From d2b7bd5328872943202d09f8dfc2f1a8214081c3 Mon Sep 17 00:00:00 2001 From: Damola09 Date: Mon, 28 Sep 2026 01:01:41 +0100 Subject: [PATCH 1/4] feat(activity): add platform activity feed endpoint (#936) GET /activity/feed returns the last 20 platform events (investment, settlement, new_listing, fully_funded), reverse-chronological with cursor pagination for older pages. Maps ActivityType enum values onto the public event type strings and caches the uncursored first page for 30s, invalidated from the trade and dividend indexers on new activity. --- prisma/schema/activity.prisma | 1 + .../activity/activity-feed.controllers.ts | 32 +++++ .../activity/activity-feed.routes.test.ts | 78 +++++++++++ src/modules/activity/activity-feed.schemas.ts | 27 ++++ src/modules/activity/activity-feed.service.ts | 130 ++++++++++++++++++ src/modules/activity/activity.routes.ts | 10 ++ .../indexer/dividend-indexer.service.ts | 6 + .../indexer/indexer-pipeline.service.ts | 5 + 8 files changed, 289 insertions(+) create mode 100644 src/modules/activity/activity-feed.controllers.ts create mode 100644 src/modules/activity/activity-feed.routes.test.ts create mode 100644 src/modules/activity/activity-feed.schemas.ts create mode 100644 src/modules/activity/activity-feed.service.ts diff --git a/prisma/schema/activity.prisma b/prisma/schema/activity.prisma index aaa9a6b4..f7c27ce2 100644 --- a/prisma/schema/activity.prisma +++ b/prisma/schema/activity.prisma @@ -17,6 +17,7 @@ enum ActivityType { GOVERNANCE_PROPOSAL_CREATED TIMELOCK_LOCKUP_PROPOSED CIRCUIT_BREAKER_THRESHOLD_UPDATED + SUPPLY_FULLY_FUNDED } model Activity { diff --git a/src/modules/activity/activity-feed.controllers.ts b/src/modules/activity/activity-feed.controllers.ts new file mode 100644 index 00000000..6164d616 --- /dev/null +++ b/src/modules/activity/activity-feed.controllers.ts @@ -0,0 +1,32 @@ +import { AsyncController } from '../../types/auth.types'; +import { ActivityFeedQuerySchema } from './activity-feed.schemas'; +import { getActivityFeed } from './activity-feed.service'; +import { + sendSuccess, + sendValidationError, +} from '../../utils/api-response.utils'; + +export const httpGetPlatformActivityFeed: AsyncController = async ( + req, + res, + next +) => { + try { + const parsed = ActivityFeedQuerySchema.safeParse(req.query); + if (!parsed.success) { + return sendValidationError( + res, + 'Invalid query parameters', + parsed.error.issues.map(issue => ({ + field: issue.path.join('.'), + message: issue.message, + })) + ); + } + + const result = await getActivityFeed(parsed.data.cursor); + sendSuccess(res, result); + } catch (error) { + next(error); + } +}; diff --git a/src/modules/activity/activity-feed.routes.test.ts b/src/modules/activity/activity-feed.routes.test.ts new file mode 100644 index 00000000..59f2d504 --- /dev/null +++ b/src/modules/activity/activity-feed.routes.test.ts @@ -0,0 +1,78 @@ +import request from 'supertest'; +import app from '../../app'; +import * as activityFeedService from './activity-feed.service'; +import { disconnectRedis } from '../../utils/redis.utils'; + +jest.mock('../../utils/redis.utils', () => { + const actual = jest.requireActual('../../utils/redis.utils'); + return { + ...actual, + cacheGetJson: jest.fn().mockResolvedValue(null), + cacheSetJson: jest.fn().mockResolvedValue(undefined), + cacheInvalidate: jest.fn().mockResolvedValue(undefined), + }; +}); + +describe('GET /activity/feed', () => { + afterAll(async () => { + await disconnectRedis(); + }); + + afterEach(() => { + jest.restoreAllMocks(); + }); + + it('returns the platform activity feed, reverse-chronological', async () => { + const mockResult: activityFeedService.ActivityFeedResult = { + items: [ + { + type: 'investment', + invoice_id: 'act-1', + amount: '100', + wallet: 'GABC…WXYZ', + timestamp: '2026-09-27T00:00:00.000Z', + }, + { + type: 'new_listing', + invoice_id: 'act-2', + amount: null, + wallet: 'GDEF…UVWX', + timestamp: '2026-09-26T00:00:00.000Z', + }, + ], + next_cursor: null, + has_more: false, + }; + + jest + .spyOn(activityFeedService, 'getActivityFeed') + .mockResolvedValueOnce(mockResult); + + const res = await request(app).get('/api/v1/activity/feed'); + expect(res.status).toBe(200); + expect(res.body.success).toBe(true); + expect(res.body.data.items).toHaveLength(2); + expect(res.body.data.items[0].type).toBe('investment'); + expect(res.body.data.has_more).toBe(false); + }); + + it('paginates with a cursor for older events', async () => { + const spy = jest + .spyOn(activityFeedService, 'getActivityFeed') + .mockResolvedValueOnce({ + items: [], + next_cursor: null, + has_more: false, + }); + + const res = await request(app).get('/api/v1/activity/feed?cursor=act-2'); + expect(res.status).toBe(200); + expect(spy).toHaveBeenCalledWith('act-2'); + }); + + it('returns 400 for unknown query parameters', async () => { + const res = await request(app).get('/api/v1/activity/feed?bogus=1'); + expect(res.status).toBe(400); + expect(res.body.success).toBe(false); + }); +}); diff --git a/src/modules/activity/activity-feed.schemas.ts b/src/modules/activity/activity-feed.schemas.ts new file mode 100644 index 00000000..abc4f9ee --- /dev/null +++ b/src/modules/activity/activity-feed.schemas.ts @@ -0,0 +1,27 @@ +import { z } from 'zod'; + +export const ActivityFeedQuerySchema = z + .object({ + cursor: z.string().optional(), + }) + .strict(); + +export type ActivityFeedQueryType = z.infer; + +export const PLATFORM_ACTIVITY_EVENT_TYPES = [ + 'investment', + 'settlement', + 'new_listing', + 'fully_funded', +] as const; + +export type PlatformActivityEventType = + (typeof PLATFORM_ACTIVITY_EVENT_TYPES)[number]; + +export interface PlatformActivityFeedItem { + type: PlatformActivityEventType; + invoice_id: string; + amount: string | null; + wallet: string; + timestamp: string; +} diff --git a/src/modules/activity/activity-feed.service.ts b/src/modules/activity/activity-feed.service.ts new file mode 100644 index 00000000..708145ca --- /dev/null +++ b/src/modules/activity/activity-feed.service.ts @@ -0,0 +1,130 @@ +// src/modules/activity/activity-feed.service.ts +// +// Platform-wide activity feed (#936): the last 20 events across the whole +// platform, reverse-chronological, with cursor pagination for older pages. +// +// The Activity model's ActivityType enum is an internal/technical taxonomy +// (KEY_BOUGHT, DIVIDEND_DISTRIBUTED, ...). The public feed exposes a small, +// stable set of event type strings instead, so this module maps enum values +// onto them. + +import { prisma } from '../../utils/prisma.utils'; +import { + cacheGetJson, + cacheSetJson, + cacheInvalidate, +} from '../../utils/redis.utils'; +import { truncateWallet } from '../../utils/wallet-display.utils'; +import { paginateQuery } from '../../utils/pagination.utils'; +import { + PlatformActivityEventType, + PlatformActivityFeedItem, +} from './activity-feed.schemas'; + +export const ACTIVITY_FEED_PAGE_SIZE = 20; +export const ACTIVITY_FEED_CACHE_KEY = 'activity:feed:v1'; +export const ACTIVITY_FEED_CACHE_TTL_SECONDS = 30; + +/** + * Maps the internal ActivityType enum onto the public-facing event type + * strings the platform feed exposes. Types not in this map are excluded from + * the feed entirely (e.g. PROFILE_UPDATED has no public-facing equivalent). + */ +const ACTIVITY_TYPE_TO_FEED_EVENT: Partial< + Record +> = { + KEY_BOUGHT: 'investment', + DIVIDEND_DISTRIBUTED: 'settlement', + CREATOR_REGISTERED: 'new_listing', + SUPPLY_FULLY_FUNDED: 'fully_funded', +}; + +const FEED_ACTIVITY_TYPES = Object.keys(ACTIVITY_TYPE_TO_FEED_EVENT); + +function mapToFeedItem(activity: { + id: string; + type: string; + actor: string; + payload: unknown; + createdAt: Date; +}): PlatformActivityFeedItem | null { + const eventType = ACTIVITY_TYPE_TO_FEED_EVENT[activity.type]; + if (!eventType) return null; + + const payload = (activity.payload as Record) || {}; + const rawAmount = + payload.amount ?? payload.price ?? payload.dividendAmount ?? null; + + return { + type: eventType, + invoice_id: activity.id, + amount: + rawAmount === null || rawAmount === undefined + ? null + : String(rawAmount), + wallet: activity.actor ? truncateWallet(activity.actor) : '', + timestamp: activity.createdAt.toISOString(), + }; +} + +export interface ActivityFeedResult { + items: PlatformActivityFeedItem[]; + next_cursor: string | null; + has_more: boolean; +} + +/** + * Fetches the platform activity feed. The uncursored first page is cached + * for 30s (ACTIVITY_FEED_CACHE_TTL_SECONDS); subsequent cursor pages read + * through to the database since they represent immutable older history. + */ +export async function getActivityFeed( + cursor?: string +): Promise { + if (!cursor) { + const cached = await cacheGetJson( + ACTIVITY_FEED_CACHE_KEY + ); + if (cached !== null) { + return cached; + } + } + + const { data, nextCursor, hasMore } = await paginateQuery( + args => + prisma.activity.findMany({ + where: { type: { in: FEED_ACTIVITY_TYPES as any } }, + orderBy: { createdAt: 'desc' }, + ...args, + }), + { + cursor: cursor ? { id: cursor } : undefined, + limit: ACTIVITY_FEED_PAGE_SIZE, + } + ); + + const items = data + .map(mapToFeedItem) + .filter((item): item is PlatformActivityFeedItem => item !== null); + + const result: ActivityFeedResult = { + items, + next_cursor: hasMore ? (nextCursor ?? null) : null, + has_more: hasMore, + }; + + if (!cursor) { + await cacheSetJson( + ACTIVITY_FEED_CACHE_KEY, + result, + ACTIVITY_FEED_CACHE_TTL_SECONDS + ); + } + + return result; +} + +/** Invalidate the cached first page. Call this whenever a new Activity row is created. */ +export async function invalidateActivityFeedCache(): Promise { + await cacheInvalidate(ACTIVITY_FEED_CACHE_KEY); +} diff --git a/src/modules/activity/activity.routes.ts b/src/modules/activity/activity.routes.ts index f5d15766..709667b1 100644 --- a/src/modules/activity/activity.routes.ts +++ b/src/modules/activity/activity.routes.ts @@ -1,10 +1,20 @@ import { Router } from 'express'; import { httpGetActivityFeed } from './activity.controllers'; +import { httpGetPlatformActivityFeed } from './activity-feed.controllers'; import { cacheControl } from '../../middlewares/cache-control.middleware'; import { ACTIVITY_FEED_CACHE_PRESET } from '../../constants/activity-feed-cache.constants'; const activityRouter = Router(); +/** + * GET /api/v1/activity/feed + * + * Platform-wide activity feed (#936): last 20 events, reverse-chronological, + * cursor pagination for older events. Must be registered before the `/` + * catch-all filtered feed below. + */ +activityRouter.get('/feed', httpGetPlatformActivityFeed); + /** * GET /api/v1/activity * diff --git a/src/modules/indexer/dividend-indexer.service.ts b/src/modules/indexer/dividend-indexer.service.ts index 375f8a75..c84cecda 100644 --- a/src/modules/indexer/dividend-indexer.service.ts +++ b/src/modules/indexer/dividend-indexer.service.ts @@ -156,6 +156,12 @@ export async function processDividendEvents( }, }); + // Invalidate the platform activity feed's cached first page (#936) so + // this new DIVIDEND_DISTRIBUTED ("settlement") activity shows up promptly. + const { invalidateActivityFeedCache } = + await import('../activity/activity-feed.service'); + await invalidateActivityFeedCache(); + logger.info( { distributionId: distribution.id, diff --git a/src/modules/indexer/indexer-pipeline.service.ts b/src/modules/indexer/indexer-pipeline.service.ts index a830328f..b03ebe07 100644 --- a/src/modules/indexer/indexer-pipeline.service.ts +++ b/src/modules/indexer/indexer-pipeline.service.ts @@ -17,6 +17,7 @@ import { logSellTransactionConfirmed } from '../../utils/sell-transaction-logger import { persistCirculatingSupply } from './persist-circulating-supply.service'; import { invalidateVolumeLeaderboardCache } from '../creators/creator-leaderboard-volume.service'; import { invalidateCreatorPortfolioStatsCache } from '../creators/creator-portfolio.service'; +import { invalidateActivityFeedCache } from '../activity/activity-feed.service'; /** * Processes a batch of on-chain trade events (KEY_BOUGHT or KEY_SOLD). @@ -97,6 +98,10 @@ export async function processTradeEvents( // instead of waiting out the full TTL (#785). await invalidateVolumeLeaderboardCache(); + // Invalidate the platform activity feed's cached first page (#936) so + // this new KEY_BOUGHT ("investment") activity shows up promptly. + await invalidateActivityFeedCache(); + // 2. Ownership read model (#897): // - buys go through recordKeyPurchase so the weighted-average cost // basis is updated on every buy (reset when rebuilding from zero). From dc947df1b192b861bea82c7a6036b4a592d5db22 Mon Sep 17 00:00:00 2001 From: Damola09 Date: Mon, 28 Sep 2026 01:02:15 +0100 Subject: [PATCH 2/4] feat(governance): track proposal quorum escalation state (#935) Adds escalationCount/originalDeadline/extendedDeadline/maxEscalations/ escalatedAt to GovernanceProposal, a governance-escalation-indexer service that syncs state from ProposalExtended contract events (with its own EventLog for idempotency) and logs a structured admin alert at max escalations, and GET /governance/proposals/escalating returning escalating proposals with participation rate (vote weight sum / totalVotingWeight). Also registers the new governance, factory, and lp routers in modules/index.ts (all three land in this branch together). --- prisma/schema/governance.prisma | 22 +++ .../governance-escalation.schemas.ts | 11 ++ .../governance-escalation.service.ts | 100 ++++++++++++ .../governance/governance.routes.test.ts | 82 ++++++++++ src/modules/governance/governance.routes.ts | 38 +++++ src/modules/index.ts | 6 + ...ernance-escalation-indexer.service.test.ts | 97 ++++++++++++ .../governance-escalation-indexer.service.ts | 146 ++++++++++++++++++ 8 files changed, 502 insertions(+) create mode 100644 src/modules/governance/governance-escalation.schemas.ts create mode 100644 src/modules/governance/governance-escalation.service.ts create mode 100644 src/modules/governance/governance.routes.test.ts create mode 100644 src/modules/governance/governance.routes.ts create mode 100644 src/modules/indexer/governance-escalation-indexer.service.test.ts create mode 100644 src/modules/indexer/governance-escalation-indexer.service.ts diff --git a/prisma/schema/governance.prisma b/prisma/schema/governance.prisma index 00f98c92..4368f53f 100644 --- a/prisma/schema/governance.prisma +++ b/prisma/schema/governance.prisma @@ -12,12 +12,34 @@ model GovernanceProposal { expiresAt DateTime closedAt DateTime? status String @default("active") // active | closed + + // Escalation tracking (#935), synced from ProposalExtended contract events. + escalationCount Int @default(0) + originalDeadline DateTime? + extendedDeadline DateTime? + maxEscalations Int @default(3) + escalatedAt DateTime? + createdAt DateTime @default(now()) updatedAt DateTime @updatedAt @@unique([keyId, proposalId]) @@index([keyId, status]) @@index([keyId, expiresAt]) + @@index([escalationCount, status]) +} + +model GovernanceEscalationEventLog { + id String @id @default(cuid()) + keyId String + proposalId String + ledger Int? + txHash String + eventIndex Int + createdAt DateTime @default(now()) + + @@unique([txHash, eventIndex]) + @@index([keyId, proposalId]) } model PauseProposal { diff --git a/src/modules/governance/governance-escalation.schemas.ts b/src/modules/governance/governance-escalation.schemas.ts new file mode 100644 index 00000000..dbed52e2 --- /dev/null +++ b/src/modules/governance/governance-escalation.schemas.ts @@ -0,0 +1,11 @@ +import { z } from 'zod'; + +export const EscalatingProposalsQuerySchema = z + .object({ + cursor: z.string().optional(), + }) + .strict(); + +export type EscalatingProposalsQueryType = z.infer< + typeof EscalatingProposalsQuerySchema +>; diff --git a/src/modules/governance/governance-escalation.service.ts b/src/modules/governance/governance-escalation.service.ts new file mode 100644 index 00000000..a6556c80 --- /dev/null +++ b/src/modules/governance/governance-escalation.service.ts @@ -0,0 +1,100 @@ +// src/modules/governance/governance-escalation.service.ts +// +// Governance proposal quorum escalation tracking (#935): proposals currently +// in escalation, with extended deadline, escalation count, and participation +// rate (sum of GovernanceVote.weight for the proposal / totalVotingWeight). + +import { prisma } from '../../utils/prisma.utils'; +import { paginateQuery } from '../../utils/pagination.utils'; + +export const ESCALATING_PROPOSALS_PAGE_SIZE = 20; + +export interface EscalatingProposalItem { + keyId: string; + proposalId: string; + title: string; + status: string; + escalationCount: number; + maxEscalations: number; + originalDeadline: string | null; + extendedDeadline: string | null; + escalatedAt: string | null; + /** + * Sum of cast vote weight divided by totalVotingWeight, as a number in + * [0, 1]. `null` when totalVotingWeight is zero (rate undefined). + */ + participationRate: number | null; +} + +async function computeParticipationRate( + keyId: string, + proposalId: string, + totalVotingWeight: string +): Promise { + const total = Number(totalVotingWeight); + if (!Number.isFinite(total) || total <= 0) { + return null; + } + + const votes = await prisma.governanceVote.aggregate({ + where: { keyId, proposalId }, + _sum: { weight: true }, + }); + + const castWeight = Number(votes._sum.weight ?? 0); + return castWeight / total; +} + +/** + * Fetches proposals currently in escalation: escalationCount > 0 and + * status = 'active', cursor-paginated, most recently escalated first. + */ +export async function getEscalatingProposals(cursor?: string): Promise<{ + items: EscalatingProposalItem[]; + next_cursor: string | null; + has_more: boolean; +}> { + const { data, nextCursor, hasMore } = await paginateQuery( + args => + prisma.governanceProposal.findMany({ + where: { escalationCount: { gt: 0 }, status: 'active' }, + orderBy: { escalatedAt: 'desc' }, + ...args, + }), + { + cursor: cursor ? { id: cursor } : undefined, + limit: ESCALATING_PROPOSALS_PAGE_SIZE, + } + ); + + const items: EscalatingProposalItem[] = await Promise.all( + data.map(async proposal => ({ + keyId: proposal.keyId, + proposalId: proposal.proposalId, + title: proposal.title, + status: proposal.status, + escalationCount: proposal.escalationCount, + maxEscalations: proposal.maxEscalations, + originalDeadline: proposal.originalDeadline + ? proposal.originalDeadline.toISOString() + : null, + extendedDeadline: proposal.extendedDeadline + ? proposal.extendedDeadline.toISOString() + : null, + escalatedAt: proposal.escalatedAt + ? proposal.escalatedAt.toISOString() + : null, + participationRate: await computeParticipationRate( + proposal.keyId, + proposal.proposalId, + proposal.totalVotingWeight + ), + })) + ); + + return { + items, + next_cursor: hasMore ? (nextCursor ?? null) : null, + has_more: hasMore, + }; +} diff --git a/src/modules/governance/governance.routes.test.ts b/src/modules/governance/governance.routes.test.ts new file mode 100644 index 00000000..e316cff6 --- /dev/null +++ b/src/modules/governance/governance.routes.test.ts @@ -0,0 +1,82 @@ +import request from 'supertest'; +import app from '../../app'; +import * as governanceEscalationService from './governance-escalation.service'; +import { disconnectRedis } from '../../utils/redis.utils'; + +jest.mock('../../utils/redis.utils', () => { + const actual = jest.requireActual('../../utils/redis.utils'); + return { + ...actual, + cacheGetJson: jest.fn().mockResolvedValue(null), + cacheSetJson: jest.fn().mockResolvedValue(undefined), + cacheInvalidate: jest.fn().mockResolvedValue(undefined), + }; +}); + +describe('GET /governance/proposals/escalating', () => { + afterAll(async () => { + await disconnectRedis(); + }); + + afterEach(() => { + jest.restoreAllMocks(); + }); + + it('returns escalating proposals with participation rate and deadlines', async () => { + jest + .spyOn(governanceEscalationService, 'getEscalatingProposals') + .mockResolvedValueOnce({ + items: [ + { + keyId: 'key-1', + proposalId: 'prop-1', + title: 'Increase treasury allocation', + status: 'active', + escalationCount: 2, + maxEscalations: 3, + originalDeadline: '2026-09-01T00:00:00.000Z', + extendedDeadline: '2026-09-15T00:00:00.000Z', + escalatedAt: '2026-09-10T00:00:00.000Z', + participationRate: 0.42, + }, + ], + next_cursor: null, + has_more: false, + }); + + const res = await request(app).get( + '/api/v1/governance/proposals/escalating' + ); + expect(res.status).toBe(200); + expect(res.body.success).toBe(true); + expect(res.body.data.items).toHaveLength(1); + expect(res.body.data.items[0].escalationCount).toBe(2); + expect(res.body.data.items[0].participationRate).toBe(0.42); + expect(res.body.data.items[0].extendedDeadline).toBe( + '2026-09-15T00:00:00.000Z' + ); + }); + + it('returns an empty list when nothing is escalating', async () => { + jest + .spyOn(governanceEscalationService, 'getEscalatingProposals') + .mockResolvedValueOnce({ + items: [], + next_cursor: null, + has_more: false, + }); + + const res = await request(app).get( + '/api/v1/governance/proposals/escalating' + ); + expect(res.status).toBe(200); + expect(res.body.data.items).toEqual([]); + }); + + it('returns 400 for unknown query parameters', async () => { + const res = await request(app).get( + '/api/v1/governance/proposals/escalating?bogus=1' + ); + expect(res.status).toBe(400); + }); +}); diff --git a/src/modules/governance/governance.routes.ts b/src/modules/governance/governance.routes.ts new file mode 100644 index 00000000..b6a565ef --- /dev/null +++ b/src/modules/governance/governance.routes.ts @@ -0,0 +1,38 @@ +import { Router } from 'express'; +import { EscalatingProposalsQuerySchema } from './governance-escalation.schemas'; +import { getEscalatingProposals } from './governance-escalation.service'; +import { + sendSuccess, + sendValidationError, +} from '../../utils/api-response.utils'; +import { zodIssuesToDetails } from '../../utils/api-response.utils'; + +const router = Router(); + +/** + * GET /governance/proposals/escalating + * + * Proposals currently in escalation (escalationCount > 0, status = active), + * showing extended deadline, escalation count, and participation rate. + * Cursor-paginated. + */ +router.get('/proposals/escalating', async (req, res, next) => { + const parsed = EscalatingProposalsQuerySchema.safeParse(req.query); + if (!parsed.success) { + sendValidationError( + res, + 'Invalid query parameters', + zodIssuesToDetails(parsed.error.issues) + ); + return; + } + + try { + const result = await getEscalatingProposals(parsed.data.cursor); + sendSuccess(res, result); + } catch (error) { + next(error); + } +}); + +export default router; diff --git a/src/modules/index.ts b/src/modules/index.ts index cb1bda93..54b2b991 100644 --- a/src/modules/index.ts +++ b/src/modules/index.ts @@ -32,6 +32,9 @@ import portfolioRouter from './portfolio/portfolio.routes'; import contractsRouter from './contracts/contract.routes'; import stakingRouter from './staking/staking.routes'; import sellersRouter from './sellers/sellers.routes'; +import governanceRouter from './governance/governance.routes'; +import factoryRouter from './factory/factory.routes'; +import lpRouter from './lp/lp.routes'; import { BASE as CREATORS_BASE } from '../constants/creator.constants'; const router = Router(); @@ -86,5 +89,8 @@ router.use('/watchlist', routeBodySizeLimit('default'), watchlistRouter); router.use('/investor/watchlist', routeBodySizeLimit('default'), watchlistRouter); router.use('/staking', routeBodySizeLimit('default'), stakingRouter); router.use('/sellers', routeBodySizeLimit('default'), sellersRouter); +router.use('/governance', routeBodySizeLimit('default'), governanceRouter); +router.use('/factory', routeBodySizeLimit('default'), factoryRouter); +router.use('/lp', routeBodySizeLimit('default'), lpRouter); export default router; diff --git a/src/modules/indexer/governance-escalation-indexer.service.test.ts b/src/modules/indexer/governance-escalation-indexer.service.test.ts new file mode 100644 index 00000000..8e7fda16 --- /dev/null +++ b/src/modules/indexer/governance-escalation-indexer.service.test.ts @@ -0,0 +1,97 @@ +// src/modules/indexer/governance-escalation-indexer.service.test.ts +import { processGovernanceEscalationEvents } from './governance-escalation-indexer.service'; +import { ProposalExtendedChainEvent } from './governance-escalation-indexer.service'; + +jest.mock('../../utils/prisma.utils', () => { + const mockTx = { + governanceEscalationEventLog: { + create: jest.fn(), + }, + governanceProposal: { + findUnique: jest.fn(), + update: jest.fn(), + }, + }; + return { + prisma: { + $transaction: jest.fn(callback => callback(mockTx)), + }, + _mockTx: mockTx, + }; +}); + +describe('processGovernanceEscalationEvents', () => { + const { _mockTx } = jest.requireMock('../../utils/prisma.utils'); + + beforeEach(() => { + jest.clearAllMocks(); + }); + + const baseEvent: ProposalExtendedChainEvent = { + eventType: 'PROPOSAL_EXTENDED', + txHash: 'tx-1', + eventIndex: 0, + ledger: 1000, + keyId: 'key-1', + proposalId: 'prop-1', + newDeadline: '2026-10-01T00:00:00.000Z', + }; + + it('increments escalationCount and updates deadline fields', async () => { + _mockTx.governanceProposal.findUnique.mockResolvedValue({ + keyId: 'key-1', + proposalId: 'prop-1', + escalationCount: 0, + maxEscalations: 3, + expiresAt: new Date('2026-09-01T00:00:00.000Z'), + originalDeadline: null, + }); + + await processGovernanceEscalationEvents([baseEvent]); + + expect(_mockTx.governanceEscalationEventLog.create).toHaveBeenCalledWith({ + data: { + keyId: 'key-1', + proposalId: 'prop-1', + ledger: 1000, + txHash: 'tx-1', + eventIndex: 0, + }, + }); + + expect(_mockTx.governanceProposal.update).toHaveBeenCalledWith({ + where: { keyId_proposalId: { keyId: 'key-1', proposalId: 'prop-1' } }, + data: expect.objectContaining({ + escalationCount: 1, + originalDeadline: new Date('2026-09-01T00:00:00.000Z'), + extendedDeadline: new Date('2026-10-01T00:00:00.000Z'), + expiresAt: new Date('2026-10-01T00:00:00.000Z'), + }), + }); + }); + + it('skips events missing required fields without throwing', async () => { + const badEvent = { + eventType: 'PROPOSAL_EXTENDED', + txHash: 'tx-2', + eventIndex: 0, + ledger: 1000, + keyId: '', + proposalId: 'prop-1', + newDeadline: 'not-a-date', + } as ProposalExtendedChainEvent; + + await expect( + processGovernanceEscalationEvents([badEvent]) + ).resolves.not.toThrow(); + expect(_mockTx.governanceProposal.update).not.toHaveBeenCalled(); + }); + + it('ignores events for unknown proposals', async () => { + _mockTx.governanceProposal.findUnique.mockResolvedValue(null); + + await processGovernanceEscalationEvents([baseEvent]); + + expect(_mockTx.governanceProposal.update).not.toHaveBeenCalled(); + }); +}); diff --git a/src/modules/indexer/governance-escalation-indexer.service.ts b/src/modules/indexer/governance-escalation-indexer.service.ts new file mode 100644 index 00000000..82bc47df --- /dev/null +++ b/src/modules/indexer/governance-escalation-indexer.service.ts @@ -0,0 +1,146 @@ +// src/modules/indexer/governance-escalation-indexer.service.ts +// +// Syncs governance proposal quorum escalation state (#935) from +// ProposalExtended contract events: increments escalationCount, moves +// expiresAt to the new (extended) deadline, and stamps escalatedAt. +// +// When escalationCount reaches maxEscalations, logs a structured admin-alert +// entry (`type: 'governance_escalation_max_reached'`) — the codebase has no +// dedicated admin-notification utility, so a structured log is the +// house-standard way to surface this for operators/alerting pipelines. + +import { Prisma } from '@prisma/client'; +import { prisma } from '../../utils/prisma.utils'; +import { logger } from '../../utils/logger.utils'; +import { buildLogFields } from '../../utils/log-fields.utils'; +import { + processIndexerChainEvents, + IndexerChainEvent, +} from '../../utils/indexer-event-processor.utils'; + +export interface ProposalExtendedChainEvent extends IndexerChainEvent { + eventType: 'PROPOSAL_EXTENDED'; + keyId: string; + proposalId: string; + /** New deadline (ISO-8601) after this escalation. */ + newDeadline: string; +} + +function isValidEvent( + event: IndexerChainEvent +): event is ProposalExtendedChainEvent { + const e = event as Partial; + return ( + event.eventType === 'PROPOSAL_EXTENDED' && + typeof e.keyId === 'string' && + e.keyId.length > 0 && + typeof e.proposalId === 'string' && + e.proposalId.length > 0 && + typeof e.newDeadline === 'string' && + !isNaN(new Date(e.newDeadline).getTime()) + ); +} + +/** + * Applies PROPOSAL_EXTENDED events to GovernanceProposal escalation fields. + * + * Each event is recorded in GovernanceEscalationEventLog inside the same + * transaction as the proposal update; a replayed event violates the + * (txHash, eventIndex) unique constraint and is skipped. + */ +export async function processGovernanceEscalationEvents( + events: IndexerChainEvent[] +): Promise { + await processIndexerChainEvents(events, async event => { + if (event.eventType !== 'PROPOSAL_EXTENDED') { + return; + } + + if (!isValidEvent(event)) { + logger.warn( + buildLogFields({ + type: 'governance_escalation_event_invalid', + eventId: `${event.txHash}:${event.eventIndex}`, + }), + 'Skipping governance escalation event with missing or invalid fields' + ); + return; + } + + const { keyId, proposalId, newDeadline } = event; + const newDeadlineDate = new Date(newDeadline); + + try { + await prisma.$transaction(async tx => { + await tx.governanceEscalationEventLog.create({ + data: { + keyId, + proposalId, + ledger: + typeof event.ledger === 'number' ? event.ledger : null, + txHash: String(event.txHash), + eventIndex: Number(event.eventIndex), + }, + }); + + const proposal = await tx.governanceProposal.findUnique({ + where: { keyId_proposalId: { keyId, proposalId } }, + }); + + if (!proposal) { + logger.warn( + buildLogFields({ + type: 'governance_escalation_proposal_not_found', + keyId, + proposalId, + eventId: `${event.txHash}:${event.eventIndex}`, + }), + 'ProposalExtended event references unknown proposal' + ); + return; + } + + const nextEscalationCount = proposal.escalationCount + 1; + const originalDeadline = + proposal.originalDeadline ?? proposal.expiresAt; + + await tx.governanceProposal.update({ + where: { keyId_proposalId: { keyId, proposalId } }, + data: { + escalationCount: nextEscalationCount, + originalDeadline, + extendedDeadline: newDeadlineDate, + expiresAt: newDeadlineDate, + escalatedAt: new Date(), + }, + }); + + if (nextEscalationCount >= proposal.maxEscalations) { + logger.warn( + buildLogFields({ + type: 'governance_escalation_max_reached', + keyId, + proposalId, + escalationCount: nextEscalationCount, + maxEscalations: proposal.maxEscalations, + extendedDeadline: newDeadlineDate, + }), + 'Governance proposal reached maximum escalations; admin attention required' + ); + } + }); + } catch (error) { + if ( + error instanceof Prisma.PrismaClientKnownRequestError && + error.code === 'P2002' + ) { + logger.debug( + { eventId: `${event.txHash}:${event.eventIndex}` }, + 'Governance escalation event already applied; skipping replay' + ); + return; + } + throw error; + } + }); +} From bddbfaa81f38413769c8ad9ecf1183590fa2cdf5 Mon Sep 17 00:00:00 2001 From: Damola09 Date: Mon, 28 Sep 2026 01:02:35 +0100 Subject: [PATCH 3/4] feat(factory): add key factory registry and indexer (#983) New FactoryDeployedKey model plus factory-indexer.service.ts syncing KeyDeployed contract events into it. GET /factory/keys?creator= lists keys deployed by a wallet in deployment order; GET /factory/keys/:contractAddress returns the factory summary (is_factory_key: true) or falls back to the general RegisteredKey lookup (is_factory_key: false), 404 only when the address isn't a key anywhere. Threads is_factory_key through the existing key summary (getCreatorProfile in creator-profile.service.ts) by checking the factory registry keyed by the profile id. --- prisma/schema/factory.prisma | 15 +++ .../creator/creator-profile.schemas.ts | 2 + .../creator/creator-profile.service.ts | 9 ++ src/modules/factory/factory.routes.test.ts | 98 +++++++++++++++++ src/modules/factory/factory.routes.ts | 62 +++++++++++ src/modules/factory/factory.schemas.ts | 11 ++ src/modules/factory/factory.service.ts | 100 ++++++++++++++++++ .../indexer/factory-indexer.service.ts | 92 ++++++++++++++++ 8 files changed, 389 insertions(+) create mode 100644 prisma/schema/factory.prisma create mode 100644 src/modules/factory/factory.routes.test.ts create mode 100644 src/modules/factory/factory.routes.ts create mode 100644 src/modules/factory/factory.schemas.ts create mode 100644 src/modules/factory/factory.service.ts create mode 100644 src/modules/indexer/factory-indexer.service.ts diff --git a/prisma/schema/factory.prisma b/prisma/schema/factory.prisma new file mode 100644 index 00000000..bd5ff280 --- /dev/null +++ b/prisma/schema/factory.prisma @@ -0,0 +1,15 @@ +// prisma/schema/factory.prisma + +model FactoryDeployedKey { + id String @id @default(cuid()) + contractAddress String @unique + creatorWallet String + keyId String? + deployedAt DateTime + txHash String + eventIndex Int + createdAt DateTime @default(now()) + + @@unique([txHash, eventIndex]) + @@index([creatorWallet]) +} diff --git a/src/modules/creator/creator-profile.schemas.ts b/src/modules/creator/creator-profile.schemas.ts index 2dbe4432..8f2791ef 100644 --- a/src/modules/creator/creator-profile.schemas.ts +++ b/src/modules/creator/creator-profile.schemas.ts @@ -60,6 +60,8 @@ export const CreatorProfileReadResponseSchema = z.object({ price24hAgo: z.string().nullable(), /** Computed percentage change. null when no baseline exists. */ priceChange24h: z.number().nullable(), + /** Whether this key was deployed through the on-chain key factory (#983). */ + is_factory_key: z.boolean(), metadata: z.object({ source: z.enum(['placeholder', 'database']), isProfileComplete: z.boolean(), diff --git a/src/modules/creator/creator-profile.service.ts b/src/modules/creator/creator-profile.service.ts index 58183329..ff625032 100644 --- a/src/modules/creator/creator-profile.service.ts +++ b/src/modules/creator/creator-profile.service.ts @@ -12,6 +12,7 @@ import { truncateString } from '../../utils/string-truncate.utils'; import { computePriceChange } from '../../utils/price-change.utils'; import { sanitizeDisplayName } from './creator-display-name-sanitize.utils'; import { invalidateKeyFeesCache } from '../keys/key-fees.service'; +import { isFactoryKey } from '../factory/factory.service'; const CREATOR_PROFILE_LIMITS = { displayName: 50, @@ -114,6 +115,7 @@ export async function getCreatorProfile( currentPrice: null, price24hAgo: null, priceChange24h: null, + is_factory_key: false, metadata: { source: 'placeholder', isProfileComplete: false, @@ -144,6 +146,12 @@ export async function getCreatorProfile( ); } + // #983: thread `is_factory_key` through the key summary. CreatorProfile + // has no dedicated on-chain contract address field, so this checks the + // factory registry keyed by the profile id, matching how `/keys/:keyId` + // already treats the id as the lookup key. + const is_factory_key = await isFactoryKey(profile.id); + return { creatorId: profile.id, displayName: profile.displayName, @@ -159,6 +167,7 @@ export async function getCreatorProfile( currentPrice: snapshot ? snapshot.currentPrice.toString() : null, price24hAgo: snapshot ? snapshot.price24hAgo.toString() : null, priceChange24h, + is_factory_key, metadata: { source: 'database', isProfileComplete: !!profile.displayName && !!profile.bio, diff --git a/src/modules/factory/factory.routes.test.ts b/src/modules/factory/factory.routes.test.ts new file mode 100644 index 00000000..812bb01c --- /dev/null +++ b/src/modules/factory/factory.routes.test.ts @@ -0,0 +1,98 @@ +import request from 'supertest'; +import app from '../../app'; +import * as factoryService from './factory.service'; +import { disconnectRedis } from '../../utils/redis.utils'; + +jest.mock('../../utils/redis.utils', () => { + const actual = jest.requireActual('../../utils/redis.utils'); + return { + ...actual, + cacheGetJson: jest.fn().mockResolvedValue(null), + cacheSetJson: jest.fn().mockResolvedValue(undefined), + cacheInvalidate: jest.fn().mockResolvedValue(undefined), + }; +}); + +describe('Factory routes', () => { + afterAll(async () => { + await disconnectRedis(); + }); + + afterEach(() => { + jest.restoreAllMocks(); + }); + + describe('GET /factory/keys', () => { + it('returns keys deployed by a creator wallet, ordered by deployment order', async () => { + jest + .spyOn(factoryService, 'getFactoryKeysByCreator') + .mockResolvedValueOnce([ + { + contractAddress: 'CADDR1', + creatorWallet: 'GCREATOR1', + keyId: 'key-1', + deployedAt: '2026-09-01T00:00:00.000Z', + is_factory_key: true, + }, + ]); + + const res = await request(app).get( + '/api/v1/factory/keys?creator=GCREATOR1' + ); + expect(res.status).toBe(200); + expect(res.body.success).toBe(true); + expect(res.body.data.items).toHaveLength(1); + expect(res.body.data.items[0].is_factory_key).toBe(true); + }); + + it('returns 400 when creator query param is missing', async () => { + const res = await request(app).get('/api/v1/factory/keys'); + expect(res.status).toBe(400); + expect(res.body.success).toBe(false); + }); + }); + + describe('GET /factory/keys/:contractAddress', () => { + it('returns is_factory_key: true for a factory-registered key', async () => { + jest + .spyOn(factoryService, 'getKeyByFactoryAddress') + .mockResolvedValueOnce({ + contractAddress: 'CADDR1', + creatorWallet: 'GCREATOR1', + keyId: 'key-1', + deployedAt: '2026-09-01T00:00:00.000Z', + is_factory_key: true, + }); + + const res = await request(app).get('/api/v1/factory/keys/CADDR1'); + expect(res.status).toBe(200); + expect(res.body.data.is_factory_key).toBe(true); + }); + + it('returns is_factory_key: false for a non-factory registered key', async () => { + jest + .spyOn(factoryService, 'getKeyByFactoryAddress') + .mockResolvedValueOnce({ + contractAddress: 'CADDR2', + creatorWallet: 'GCREATOR2', + is_factory_key: false, + }); + + const res = await request(app).get('/api/v1/factory/keys/CADDR2'); + expect(res.status).toBe(200); + expect(res.body.data.is_factory_key).toBe(false); + }); + + it('returns 404 when the address is not a key anywhere', async () => { + jest + .spyOn(factoryService, 'getKeyByFactoryAddress') + .mockRejectedValueOnce( + new factoryService.KeyAddressNotFoundError('CUNKNOWN') + ); + + const res = await request(app).get('/api/v1/factory/keys/CUNKNOWN'); + expect(res.status).toBe(404); + expect(res.body.success).toBe(false); + }); + }); +}); diff --git a/src/modules/factory/factory.routes.ts b/src/modules/factory/factory.routes.ts new file mode 100644 index 00000000..680c8c2f --- /dev/null +++ b/src/modules/factory/factory.routes.ts @@ -0,0 +1,62 @@ +import { Router } from 'express'; +import { + sendSuccess, + sendNotFound, + sendValidationError, + zodIssuesToDetails, +} from '../../utils/api-response.utils'; +import { FactoryKeysByCreatorQuerySchema } from './factory.schemas'; +import { + getFactoryKeysByCreator, + getKeyByFactoryAddress, + KeyAddressNotFoundError, +} from './factory.service'; + +const router = Router(); + +/** + * GET /factory/keys?creator= + * All keys deployed by that creator wallet, ordered by deployment order. + */ +router.get('/keys', async (req, res, next) => { + const parsed = FactoryKeysByCreatorQuerySchema.safeParse(req.query); + if (!parsed.success) { + sendValidationError( + res, + 'Invalid query parameters', + zodIssuesToDetails(parsed.error.issues) + ); + return; + } + + try { + const keys = await getFactoryKeysByCreator(parsed.data.creator); + sendSuccess(res, { items: keys }); + } catch (error) { + next(error); + } +}); + +/** + * GET /factory/keys/:contractAddress + * Looks up a key by contract address in the factory registry; falls back to + * the general key registry (is_factory_key: false) when it's a real key that + * wasn't deployed through the factory. 404 only when the address isn't a key + * anywhere. + */ +router.get('/keys/:contractAddress', async (req, res, next) => { + try { + const result = await getKeyByFactoryAddress( + String(req.params.contractAddress) + ); + sendSuccess(res, result); + } catch (error) { + if (error instanceof KeyAddressNotFoundError) { + sendNotFound(res, 'Key'); + return; + } + next(error); + } +}); + +export default router; diff --git a/src/modules/factory/factory.schemas.ts b/src/modules/factory/factory.schemas.ts new file mode 100644 index 00000000..b145a68c --- /dev/null +++ b/src/modules/factory/factory.schemas.ts @@ -0,0 +1,11 @@ +import { z } from 'zod'; + +export const FactoryKeysByCreatorQuerySchema = z + .object({ + creator: z.string().min(1, 'creator is required'), + }) + .strict(); + +export type FactoryKeysByCreatorQueryType = z.infer< + typeof FactoryKeysByCreatorQuerySchema +>; diff --git a/src/modules/factory/factory.service.ts b/src/modules/factory/factory.service.ts new file mode 100644 index 00000000..454ac19a --- /dev/null +++ b/src/modules/factory/factory.service.ts @@ -0,0 +1,100 @@ +// src/modules/factory/factory.service.ts +// +// Key factory registry API (#983): lookups over FactoryDeployedKey, plus a +// fallback to the general key registry (RegisteredKey) when an address is a +// real key but wasn't deployed through the factory. + +import { prisma } from '../../utils/prisma.utils'; + +export interface FactoryKeySummary { + contractAddress: string; + creatorWallet: string; + keyId: string | null; + deployedAt: string; + is_factory_key: true; +} + +export interface NonFactoryKeySummary { + contractAddress: string; + creatorWallet: string; + is_factory_key: false; +} + +export class KeyAddressNotFoundError extends Error { + constructor(address: string) { + super(`No key found for address: ${address}`); + this.name = 'KeyAddressNotFoundError'; + } +} + +/** + * All keys deployed by a given creator wallet, ordered by deployment order + * (deployedAt ascending). + */ +export async function getFactoryKeysByCreator( + creatorWallet: string +): Promise { + const rows = await prisma.factoryDeployedKey.findMany({ + where: { creatorWallet }, + orderBy: { deployedAt: 'asc' }, + }); + + return rows.map(row => ({ + contractAddress: row.contractAddress, + creatorWallet: row.creatorWallet, + keyId: row.keyId, + deployedAt: row.deployedAt.toISOString(), + is_factory_key: true as const, + })); +} + +/** + * Looks up a key by contract address. If it's in the factory registry, + * returns its factory summary with is_factory_key: true. Otherwise falls + * back to the general registered-key lookup and returns is_factory_key: + * false. Throws KeyAddressNotFoundError only when the address isn't a key + * anywhere. + */ +export async function getKeyByFactoryAddress( + contractAddress: string +): Promise { + const factoryKey = await prisma.factoryDeployedKey.findUnique({ + where: { contractAddress }, + }); + + if (factoryKey) { + return { + contractAddress: factoryKey.contractAddress, + creatorWallet: factoryKey.creatorWallet, + keyId: factoryKey.keyId, + deployedAt: factoryKey.deployedAt.toISOString(), + is_factory_key: true, + }; + } + + const registeredKey = await prisma.registeredKey.findUnique({ + where: { keyAddress: contractAddress }, + }); + + if (registeredKey) { + return { + contractAddress: registeredKey.keyAddress, + creatorWallet: registeredKey.creatorWallet, + is_factory_key: false, + }; + } + + throw new KeyAddressNotFoundError(contractAddress); +} + +/** + * Checks whether a contract address exists in the factory registry. + * Used to thread `is_factory_key` through existing key summary responses. + */ +export async function isFactoryKey(contractAddress: string): Promise { + const factoryKey = await prisma.factoryDeployedKey.findUnique({ + where: { contractAddress }, + select: { id: true }, + }); + return factoryKey !== null; +} diff --git a/src/modules/indexer/factory-indexer.service.ts b/src/modules/indexer/factory-indexer.service.ts new file mode 100644 index 00000000..1ffb5579 --- /dev/null +++ b/src/modules/indexer/factory-indexer.service.ts @@ -0,0 +1,92 @@ +// src/modules/indexer/factory-indexer.service.ts +// +// Key factory deployment event indexing (#983): KeyDeployed contract events +// insert a FactoryDeployedKey row into the factory registry. contractAddress +// is itself unique (a re-deploy replay is naturally idempotent), but the +// standard (txHash, eventIndex) unique-constraint + P2002-catch pattern is +// still followed for consistency with the rest of the indexer suite. + +import { Prisma } from '@prisma/client'; +import { prisma } from '../../utils/prisma.utils'; +import { logger } from '../../utils/logger.utils'; +import { buildLogFields } from '../../utils/log-fields.utils'; +import { + processIndexerChainEvents, + IndexerChainEvent, +} from '../../utils/indexer-event-processor.utils'; + +export interface KeyDeployedChainEvent extends IndexerChainEvent { + eventType: 'KEY_DEPLOYED'; + contractAddress: string; + creatorWallet: string; + keyId?: string; + deployedAt: string; +} + +function isValidEvent( + event: IndexerChainEvent +): event is KeyDeployedChainEvent { + const e = event as Partial; + return ( + event.eventType === 'KEY_DEPLOYED' && + typeof e.contractAddress === 'string' && + e.contractAddress.length > 0 && + typeof e.creatorWallet === 'string' && + e.creatorWallet.length > 0 && + typeof e.deployedAt === 'string' && + !isNaN(new Date(e.deployedAt).getTime()) + ); +} + +/** + * Applies KEY_DEPLOYED events to the FactoryDeployedKey registry. + */ +export async function processFactoryEvents( + events: IndexerChainEvent[] +): Promise { + await processIndexerChainEvents(events, async event => { + if (event.eventType !== 'KEY_DEPLOYED') { + return; + } + + if (!isValidEvent(event)) { + logger.warn( + buildLogFields({ + type: 'factory_event_invalid', + eventId: `${event.txHash}:${event.eventIndex}`, + }), + 'Skipping factory event with missing or invalid fields' + ); + return; + } + + const { contractAddress, creatorWallet, keyId, deployedAt } = event; + + try { + await prisma.$transaction(async tx => { + await tx.factoryDeployedKey.create({ + data: { + contractAddress, + creatorWallet, + keyId: keyId ?? null, + deployedAt: new Date(deployedAt), + txHash: String(event.txHash), + eventIndex: Number(event.eventIndex), + }, + }); + }); + } catch (error) { + if ( + error instanceof Prisma.PrismaClientKnownRequestError && + error.code === 'P2002' + ) { + logger.debug( + { eventId: `${event.txHash}:${event.eventIndex}` }, + 'Factory deployment event already applied; skipping replay' + ); + return; + } + throw error; + } + }); +} From ec9a3de70a058026b0c5f9440b3dbca076158881 Mon Sep 17 00:00:00 2001 From: Damola09 Date: Mon, 28 Sep 2026 01:02:56 +0100 Subject: [PATCH 4/4] feat(lp): add LP position indexing and rewards API (#980) New LpPosition/LpEventLog models plus lp-indexer.service.ts handling LP_ADDED/LP_CLAIMED/LP_REMOVED events. accrueLpRewards pro-rates a simplified reward pool across active positions by sharePercent, called from trade-indexer.service.ts after each processed trade. GET /lp/positions?wallet= and GET /lp/positions/:lpId are JWT-auth'd and wallet-scoped (404 on any not-owned/not-found lpId, no existence leakage); GET /lp/pool/:keyId is public and returns pool size plus a clearly-labelled simplified APR estimate. Also adds the combined migration for all four issues' schema changes (activity SUPPLY_FULLY_FUNDED enum value, governance escalation fields + GovernanceEscalationEventLog, FactoryDeployedKey, LpPosition/ LpEventLog). --- prisma/schema/liquidity.prisma | 26 ++ .../migration.sql | 92 +++++++ .../indexer/lp-indexer.service.test.ts | 133 ++++++++++ src/modules/indexer/lp-indexer.service.ts | 227 ++++++++++++++++++ src/modules/indexer/trade-indexer.service.ts | 7 + src/modules/lp/lp.routes.test.ts | 120 +++++++++ src/modules/lp/lp.routes.ts | 95 ++++++++ src/modules/lp/lp.schemas.ts | 11 + src/modules/lp/lp.service.ts | 114 +++++++++ 9 files changed, 825 insertions(+) create mode 100644 prisma/schema/liquidity.prisma create mode 100644 prisma/schema/migrations/20260927000000_platform_endpoints_936_935_983_980/migration.sql create mode 100644 src/modules/indexer/lp-indexer.service.test.ts create mode 100644 src/modules/indexer/lp-indexer.service.ts create mode 100644 src/modules/lp/lp.routes.test.ts create mode 100644 src/modules/lp/lp.routes.ts create mode 100644 src/modules/lp/lp.schemas.ts create mode 100644 src/modules/lp/lp.service.ts diff --git a/prisma/schema/liquidity.prisma b/prisma/schema/liquidity.prisma new file mode 100644 index 00000000..b880d4b0 --- /dev/null +++ b/prisma/schema/liquidity.prisma @@ -0,0 +1,26 @@ +// prisma/schema/liquidity.prisma + +model LpPosition { + id String @id @default(cuid()) + lpId String @unique @default(cuid()) + wallet String + keyId String + sharePercent Decimal @default(0) + accruedRewards Decimal @default(0) + status String @default("active") // active | removed + createdAt DateTime @default(now()) + updatedAt DateTime @updatedAt + + @@index([wallet]) + @@index([keyId]) +} + +model LpEventLog { + id String @id @default(cuid()) + txHash String + eventIndex Int + eventType String + createdAt DateTime @default(now()) + + @@unique([txHash, eventIndex]) +} diff --git a/prisma/schema/migrations/20260927000000_platform_endpoints_936_935_983_980/migration.sql b/prisma/schema/migrations/20260927000000_platform_endpoints_936_935_983_980/migration.sql new file mode 100644 index 00000000..de81d2b6 --- /dev/null +++ b/prisma/schema/migrations/20260927000000_platform_endpoints_936_935_983_980/migration.sql @@ -0,0 +1,92 @@ +-- #936: platform activity feed — new ActivityType enum value +ALTER TYPE "ActivityType" ADD VALUE 'SUPPLY_FULLY_FUNDED'; + +-- #935: governance proposal quorum escalation tracking +ALTER TABLE "GovernanceProposal" ADD COLUMN "escalationCount" INTEGER NOT NULL DEFAULT 0, +ADD COLUMN "originalDeadline" TIMESTAMP(3), +ADD COLUMN "extendedDeadline" TIMESTAMP(3), +ADD COLUMN "maxEscalations" INTEGER NOT NULL DEFAULT 3, +ADD COLUMN "escalatedAt" TIMESTAMP(3); + +-- CreateIndex +CREATE INDEX "GovernanceProposal_escalationCount_status_idx" ON "GovernanceProposal"("escalationCount", "status"); + +-- CreateTable +CREATE TABLE "GovernanceEscalationEventLog" ( + "id" TEXT NOT NULL, + "keyId" TEXT NOT NULL, + "proposalId" TEXT NOT NULL, + "ledger" INTEGER, + "txHash" TEXT NOT NULL, + "eventIndex" INTEGER NOT NULL, + "createdAt" TIMESTAMP(3) NOT NULL DEFAULT CURRENT_TIMESTAMP, + + CONSTRAINT "GovernanceEscalationEventLog_pkey" PRIMARY KEY ("id") +); + +-- CreateIndex +CREATE UNIQUE INDEX "GovernanceEscalationEventLog_txHash_eventIndex_key" ON "GovernanceEscalationEventLog"("txHash", "eventIndex"); + +-- CreateIndex +CREATE INDEX "GovernanceEscalationEventLog_keyId_proposalId_idx" ON "GovernanceEscalationEventLog"("keyId", "proposalId"); + +-- #983: key factory deployment event indexing and registry +CREATE TABLE "FactoryDeployedKey" ( + "id" TEXT NOT NULL, + "contractAddress" TEXT NOT NULL, + "creatorWallet" TEXT NOT NULL, + "keyId" TEXT, + "deployedAt" TIMESTAMP(3) NOT NULL, + "txHash" TEXT NOT NULL, + "eventIndex" INTEGER NOT NULL, + "createdAt" TIMESTAMP(3) NOT NULL DEFAULT CURRENT_TIMESTAMP, + + CONSTRAINT "FactoryDeployedKey_pkey" PRIMARY KEY ("id") +); + +-- CreateIndex +CREATE UNIQUE INDEX "FactoryDeployedKey_contractAddress_key" ON "FactoryDeployedKey"("contractAddress"); + +-- CreateIndex +CREATE UNIQUE INDEX "FactoryDeployedKey_txHash_eventIndex_key" ON "FactoryDeployedKey"("txHash", "eventIndex"); + +-- CreateIndex +CREATE INDEX "FactoryDeployedKey_creatorWallet_idx" ON "FactoryDeployedKey"("creatorWallet"); + +-- #980: LP position indexing and rewards +CREATE TABLE "LpPosition" ( + "id" TEXT NOT NULL, + "lpId" TEXT NOT NULL, + "wallet" TEXT NOT NULL, + "keyId" TEXT NOT NULL, + "sharePercent" DECIMAL(65,30) NOT NULL DEFAULT 0, + "accruedRewards" DECIMAL(65,30) NOT NULL DEFAULT 0, + "status" TEXT NOT NULL DEFAULT 'active', + "createdAt" TIMESTAMP(3) NOT NULL DEFAULT CURRENT_TIMESTAMP, + "updatedAt" TIMESTAMP(3) NOT NULL, + + CONSTRAINT "LpPosition_pkey" PRIMARY KEY ("id") +); + +-- CreateIndex +CREATE UNIQUE INDEX "LpPosition_lpId_key" ON "LpPosition"("lpId"); + +-- CreateIndex +CREATE INDEX "LpPosition_wallet_idx" ON "LpPosition"("wallet"); + +-- CreateIndex +CREATE INDEX "LpPosition_keyId_idx" ON "LpPosition"("keyId"); + +-- CreateTable +CREATE TABLE "LpEventLog" ( + "id" TEXT NOT NULL, + "txHash" TEXT NOT NULL, + "eventIndex" INTEGER NOT NULL, + "eventType" TEXT NOT NULL, + "createdAt" TIMESTAMP(3) NOT NULL DEFAULT CURRENT_TIMESTAMP, + + CONSTRAINT "LpEventLog_pkey" PRIMARY KEY ("id") +); + +-- CreateIndex +CREATE UNIQUE INDEX "LpEventLog_txHash_eventIndex_key" ON "LpEventLog"("txHash", "eventIndex"); diff --git a/src/modules/indexer/lp-indexer.service.test.ts b/src/modules/indexer/lp-indexer.service.test.ts new file mode 100644 index 00000000..6fc4fc06 --- /dev/null +++ b/src/modules/indexer/lp-indexer.service.test.ts @@ -0,0 +1,133 @@ +// src/modules/indexer/lp-indexer.service.test.ts +import { processLpEvents, LpChainEvent } from './lp-indexer.service'; + +jest.mock('../../utils/prisma.utils', () => { + const mockTx = { + lpEventLog: { create: jest.fn() }, + lpPosition: { + findFirst: jest.fn(), + create: jest.fn(), + update: jest.fn(), + }, + }; + return { + prisma: { + $transaction: jest.fn(callback => callback(mockTx)), + lpPosition: { + findMany: jest.fn(), + update: jest.fn(), + }, + }, + _mockTx: mockTx, + }; +}); + +describe('processLpEvents', () => { + const { _mockTx } = jest.requireMock('../../utils/prisma.utils'); + + beforeEach(() => { + jest.clearAllMocks(); + }); + + it('creates a new LpPosition on LP_ADDED when none exists', async () => { + _mockTx.lpPosition.findFirst.mockResolvedValue(null); + + const event: LpChainEvent = { + eventType: 'LP_ADDED', + txHash: 'tx-1', + eventIndex: 0, + ledger: 100, + wallet: 'wallet-1', + keyId: 'key-1', + sharePercent: '10', + }; + + await processLpEvents([event]); + + expect(_mockTx.lpPosition.create).toHaveBeenCalledWith({ + data: { + wallet: 'wallet-1', + keyId: 'key-1', + sharePercent: '10', + status: 'active', + }, + }); + }); + + it('updates sharePercent on LP_ADDED when a position already exists', async () => { + _mockTx.lpPosition.findFirst.mockResolvedValue({ id: 'pos-1' }); + + const event: LpChainEvent = { + eventType: 'LP_ADDED', + txHash: 'tx-2', + eventIndex: 0, + ledger: 100, + wallet: 'wallet-1', + keyId: 'key-1', + sharePercent: '20', + }; + + await processLpEvents([event]); + + expect(_mockTx.lpPosition.update).toHaveBeenCalledWith({ + where: { id: 'pos-1' }, + data: { sharePercent: '20', status: 'active' }, + }); + }); + + it('increments accruedRewards on LP_CLAIMED', async () => { + _mockTx.lpPosition.findFirst.mockResolvedValue({ id: 'pos-1' }); + + const event: LpChainEvent = { + eventType: 'LP_CLAIMED', + txHash: 'tx-3', + eventIndex: 0, + ledger: 100, + wallet: 'wallet-1', + keyId: 'key-1', + rewardAmount: '5', + }; + + await processLpEvents([event]); + + expect(_mockTx.lpPosition.update).toHaveBeenCalledWith({ + where: { id: 'pos-1' }, + data: { accruedRewards: { increment: '5' } }, + }); + }); + + it('sets status to removed on LP_REMOVED', async () => { + _mockTx.lpPosition.findFirst.mockResolvedValue({ id: 'pos-1' }); + + const event: LpChainEvent = { + eventType: 'LP_REMOVED', + txHash: 'tx-4', + eventIndex: 0, + ledger: 100, + wallet: 'wallet-1', + keyId: 'key-1', + }; + + await processLpEvents([event]); + + expect(_mockTx.lpPosition.update).toHaveBeenCalledWith({ + where: { id: 'pos-1' }, + data: { status: 'removed' }, + }); + }); + + it('skips invalid events without throwing', async () => { + const badEvent = { + eventType: 'LP_ADDED', + txHash: 'tx-5', + eventIndex: 0, + ledger: 100, + wallet: '', + keyId: 'key-1', + sharePercent: '10', + } as LpChainEvent; + + await expect(processLpEvents([badEvent])).resolves.not.toThrow(); + expect(_mockTx.lpPosition.create).not.toHaveBeenCalled(); + }); +}); diff --git a/src/modules/indexer/lp-indexer.service.ts b/src/modules/indexer/lp-indexer.service.ts new file mode 100644 index 00000000..0e5b2922 --- /dev/null +++ b/src/modules/indexer/lp-indexer.service.ts @@ -0,0 +1,227 @@ +// src/modules/indexer/lp-indexer.service.ts +// +// LP position indexing (#980): handles LP_ADDED / LP_CLAIMED / LP_REMOVED +// events. LP_ADDED creates/upserts an LpPosition with the given share. +// LP_CLAIMED adds to accruedRewards. LP_REMOVED sets status = 'removed'. + +import { Prisma } from '@prisma/client'; +import { prisma } from '../../utils/prisma.utils'; +import { logger } from '../../utils/logger.utils'; +import { buildLogFields } from '../../utils/log-fields.utils'; +import { + processIndexerChainEvents, + IndexerChainEvent, +} from '../../utils/indexer-event-processor.utils'; + +export interface LpChainEvent extends IndexerChainEvent { + eventType: 'LP_ADDED' | 'LP_CLAIMED' | 'LP_REMOVED'; + wallet: string; + keyId: string; + /** Share percent, only meaningful for LP_ADDED (0-100). */ + sharePercent?: string; + /** Reward amount claimed, only meaningful for LP_CLAIMED. */ + rewardAmount?: string; +} + +function isValidEvent(event: IndexerChainEvent): event is LpChainEvent { + const e = event as Partial; + if ( + (event.eventType !== 'LP_ADDED' && + event.eventType !== 'LP_CLAIMED' && + event.eventType !== 'LP_REMOVED') || + typeof e.wallet !== 'string' || + e.wallet.length === 0 || + typeof e.keyId !== 'string' || + e.keyId.length === 0 + ) { + return false; + } + + if (event.eventType === 'LP_ADDED') { + const share = Number(e.sharePercent); + return ( + typeof e.sharePercent === 'string' && + Number.isFinite(share) && + share >= 0 + ); + } + + if (event.eventType === 'LP_CLAIMED') { + const reward = Number(e.rewardAmount); + return ( + typeof e.rewardAmount === 'string' && + Number.isFinite(reward) && + reward >= 0 + ); + } + + return true; +} + +/** + * Applies LP_ADDED / LP_CLAIMED / LP_REMOVED events to LpPosition. + * + * Each event is recorded in LpEventLog inside the same transaction as the + * position change; a replayed event violates the (txHash, eventIndex) + * unique constraint and is skipped. + */ +export async function processLpEvents( + events: IndexerChainEvent[] +): Promise { + await processIndexerChainEvents(events, async event => { + if ( + event.eventType !== 'LP_ADDED' && + event.eventType !== 'LP_CLAIMED' && + event.eventType !== 'LP_REMOVED' + ) { + return; + } + + if (!isValidEvent(event)) { + logger.warn( + buildLogFields({ + type: 'lp_event_invalid', + eventId: `${event.txHash}:${event.eventIndex}`, + }), + 'Skipping LP event with missing or invalid fields' + ); + return; + } + + const { wallet, keyId } = event; + + try { + await prisma.$transaction(async tx => { + await tx.lpEventLog.create({ + data: { + txHash: String(event.txHash), + eventIndex: Number(event.eventIndex), + eventType: event.eventType, + }, + }); + + if (event.eventType === 'LP_ADDED') { + const sharePercent = event.sharePercent as string; + const existing = await tx.lpPosition.findFirst({ + where: { wallet, keyId }, + }); + + if (existing) { + await tx.lpPosition.update({ + where: { id: existing.id }, + data: { sharePercent, status: 'active' }, + }); + } else { + await tx.lpPosition.create({ + data: { wallet, keyId, sharePercent, status: 'active' }, + }); + } + return; + } + + if (event.eventType === 'LP_CLAIMED') { + const rewardAmount = event.rewardAmount as string; + const existing = await tx.lpPosition.findFirst({ + where: { wallet, keyId }, + }); + if (!existing) { + logger.warn( + buildLogFields({ + type: 'lp_claim_no_position', + wallet, + keyId, + eventId: `${event.txHash}:${event.eventIndex}`, + }), + 'LP_CLAIMED event references a wallet/key with no LpPosition' + ); + return; + } + await tx.lpPosition.update({ + where: { id: existing.id }, + data: { + accruedRewards: { + increment: rewardAmount, + }, + }, + }); + return; + } + + // LP_REMOVED + const existing = await tx.lpPosition.findFirst({ + where: { wallet, keyId }, + }); + if (!existing) { + logger.warn( + buildLogFields({ + type: 'lp_removed_no_position', + wallet, + keyId, + eventId: `${event.txHash}:${event.eventIndex}`, + }), + 'LP_REMOVED event references a wallet/key with no LpPosition' + ); + return; + } + await tx.lpPosition.update({ + where: { id: existing.id }, + data: { status: 'removed' }, + }); + }); + } catch (error) { + if ( + error instanceof Prisma.PrismaClientKnownRequestError && + error.code === 'P2002' + ) { + logger.debug( + { eventId: `${event.txHash}:${event.eventIndex}` }, + 'LP event already applied; skipping replay' + ); + return; + } + throw error; + } + }); +} + +/** + * Pro-rates a reward pool across active LpPosition rows for a key by + * sharePercent, called from the trade indexer after each processed trade + * (#980). This is a simplified accrual model: the reward pool per trade is a + * fixed small fraction of the trade amount, not a precise on-chain fee split + * — see the inline comment at the call site in trade-indexer.service.ts. + */ +export async function accrueLpRewards( + keyId: string, + tradeAmount: number +): Promise { + if (!keyId || !Number.isFinite(tradeAmount) || tradeAmount <= 0) { + return; + } + + const LP_REWARD_POOL_BPS = 50; // 0.5% of trade amount, simplified placeholder + const rewardPool = (tradeAmount * LP_REWARD_POOL_BPS) / 10_000; + if (rewardPool <= 0) { + return; + } + + const activePositions = await prisma.lpPosition.findMany({ + where: { keyId, status: 'active' }, + }); + + if (activePositions.length === 0) { + return; + } + + await Promise.all( + activePositions.map(position => { + const share = Number(position.sharePercent) / 100; + const reward = rewardPool * share; + if (reward <= 0) return Promise.resolve(); + return prisma.lpPosition.update({ + where: { id: position.id }, + data: { accruedRewards: { increment: reward } }, + }); + }) + ); +} diff --git a/src/modules/indexer/trade-indexer.service.ts b/src/modules/indexer/trade-indexer.service.ts index 055c61a9..15bf0ce1 100644 --- a/src/modules/indexer/trade-indexer.service.ts +++ b/src/modules/indexer/trade-indexer.service.ts @@ -87,6 +87,13 @@ export async function processTradeEvent( }, }); + try { + const { accrueLpRewards } = await import('./lp-indexer.service'); + await accrueLpRewards(event.creator_id, Number(event.price)); + } catch { + // Non-critical: LP reward accrual failure shouldn't fail trade indexing + } + try { const { invalidateCreatorDashboardCache } = await import('../creator/creator-dashboard.service'); diff --git a/src/modules/lp/lp.routes.test.ts b/src/modules/lp/lp.routes.test.ts new file mode 100644 index 00000000..178fd744 --- /dev/null +++ b/src/modules/lp/lp.routes.test.ts @@ -0,0 +1,120 @@ +import request from 'supertest'; +import app from '../../app'; +import * as lpService from './lp.service'; +import { signWalletAccessToken } from '../../utils/jwt.utils'; +import { disconnectRedis } from '../../utils/redis.utils'; + +jest.mock('../../utils/redis.utils', () => { + const actual = jest.requireActual('../../utils/redis.utils'); + return { + ...actual, + cacheGetJson: jest.fn().mockResolvedValue(null), + cacheSetJson: jest.fn().mockResolvedValue(undefined), + cacheInvalidate: jest.fn().mockResolvedValue(undefined), + }; +}); + +describe('LP routes', () => { + const wallet = 'GLPWALLET0000000000000000000000000000000000000000000'; + const token = signWalletAccessToken(wallet); + + afterAll(async () => { + await disconnectRedis(); + }); + + afterEach(() => { + jest.restoreAllMocks(); + }); + + describe('GET /api/v1/lp/positions', () => { + it('returns 401 without a token', async () => { + const res = await request(app).get( + `/api/v1/lp/positions?wallet=${wallet}` + ); + expect(res.status).toBe(401); + }); + + it('returns 403 when the query wallet does not match the authenticated wallet', async () => { + const res = await request(app) + .get('/api/v1/lp/positions?wallet=GDIFFERENTWALLET') + .set('Authorization', `Bearer ${token}`); + expect(res.status).toBe(403); + }); + + it('returns active LP positions for the authenticated wallet', async () => { + jest.spyOn(lpService, 'getLpPositionsByWallet').mockResolvedValueOnce([ + { + lpId: 'lp-1', + wallet, + keyId: 'key-1', + sharePercent: '10', + accruedRewards: '5.5', + status: 'active', + createdAt: '2026-09-01T00:00:00.000Z', + updatedAt: '2026-09-01T00:00:00.000Z', + }, + ]); + + const res = await request(app) + .get(`/api/v1/lp/positions?wallet=${wallet}`) + .set('Authorization', `Bearer ${token}`); + expect(res.status).toBe(200); + expect(res.body.data.items).toHaveLength(1); + expect(res.body.data.items[0].sharePercent).toBe('10'); + }); + }); + + describe('GET /api/v1/lp/positions/:lpId', () => { + it('returns 401 without a token', async () => { + const res = await request(app).get('/api/v1/lp/positions/lp-1'); + expect(res.status).toBe(401); + }); + + it('returns 404 when the position does not exist or is not owned by the wallet', async () => { + jest + .spyOn(lpService, 'getLpPositionById') + .mockRejectedValueOnce( + new lpService.LpPositionNotFoundError('lp-1') + ); + + const res = await request(app) + .get('/api/v1/lp/positions/lp-1') + .set('Authorization', `Bearer ${token}`); + expect(res.status).toBe(404); + }); + + it('returns position detail when found and owned', async () => { + jest.spyOn(lpService, 'getLpPositionById').mockResolvedValueOnce({ + lpId: 'lp-1', + wallet, + keyId: 'key-1', + sharePercent: '10', + accruedRewards: '5.5', + status: 'active', + createdAt: '2026-09-01T00:00:00.000Z', + updatedAt: '2026-09-01T00:00:00.000Z', + }); + + const res = await request(app) + .get('/api/v1/lp/positions/lp-1') + .set('Authorization', `Bearer ${token}`); + expect(res.status).toBe(200); + expect(res.body.data.lpId).toBe('lp-1'); + }); + }); + + describe('GET /api/v1/lp/pool/:keyId', () => { + it('returns pool size and APR estimate without auth', async () => { + jest.spyOn(lpService, 'getLpPoolSummary').mockResolvedValueOnce({ + keyId: 'key-1', + totalPoolSize: '100', + estimatedApr: 12.5, + }); + + const res = await request(app).get('/api/v1/lp/pool/key-1'); + expect(res.status).toBe(200); + expect(res.body.data.totalPoolSize).toBe('100'); + expect(res.body.data.estimatedApr).toBe(12.5); + }); + }); +}); diff --git a/src/modules/lp/lp.routes.ts b/src/modules/lp/lp.routes.ts new file mode 100644 index 00000000..426125ce --- /dev/null +++ b/src/modules/lp/lp.routes.ts @@ -0,0 +1,95 @@ +import { Router } from 'express'; +import { + sendSuccess, + sendNotFound, + sendForbidden, + sendValidationError, + zodIssuesToDetails, +} from '../../utils/api-response.utils'; +import { + requireJwtAuth, + AuthenticatedRequest, +} from '../../middlewares/jwt-auth.middleware'; +import { LpPositionsByWalletQuerySchema } from './lp.schemas'; +import { + getLpPositionsByWallet, + getLpPositionById, + getLpPoolSummary, + LpPositionNotFoundError, +} from './lp.service'; + +const router = Router(); + +/** + * GET /lp/positions?wallet= + * Requires JWT auth; the authenticated wallet must match the query wallet. + * All active LP positions for that wallet with sharePercent and accruedRewards. + */ +router.get( + '/positions', + requireJwtAuth, + async (req: AuthenticatedRequest, res, next) => { + const parsed = LpPositionsByWalletQuerySchema.safeParse(req.query); + if (!parsed.success) { + sendValidationError( + res, + 'Invalid query parameters', + zodIssuesToDetails(parsed.error.issues) + ); + return; + } + + if (req.user!.wallet !== parsed.data.wallet) { + sendForbidden(res, 'Wallet does not match the authenticated wallet'); + return; + } + + try { + const positions = await getLpPositionsByWallet(parsed.data.wallet); + sendSuccess(res, { items: positions }); + } catch (error) { + next(error); + } + } +); + +/** + * GET /lp/positions/:lpId + * Requires JWT auth; 404 both when the position doesn't exist and when it + * isn't owned by the authenticated wallet (no existence leakage). + */ +router.get( + '/positions/:lpId', + requireJwtAuth, + async (req: AuthenticatedRequest, res, next) => { + try { + const position = await getLpPositionById( + String(req.params.lpId), + req.user!.wallet + ); + sendSuccess(res, position); + } catch (error) { + if (error instanceof LpPositionNotFoundError) { + sendNotFound(res, 'LP position'); + return; + } + next(error); + } + } +); + +/** + * GET /lp/pool/:keyId + * Public (aggregate data, no auth required). Total pool size and a + * simplified APR estimate. + */ +router.get('/pool/:keyId', async (req, res, next) => { + try { + const pool = await getLpPoolSummary(String(req.params.keyId)); + sendSuccess(res, pool); + } catch (error) { + next(error); + } +}); + +export default router; diff --git a/src/modules/lp/lp.schemas.ts b/src/modules/lp/lp.schemas.ts new file mode 100644 index 00000000..e6bc3381 --- /dev/null +++ b/src/modules/lp/lp.schemas.ts @@ -0,0 +1,11 @@ +import { z } from 'zod'; + +export const LpPositionsByWalletQuerySchema = z + .object({ + wallet: z.string().min(1, 'wallet is required'), + }) + .strict(); + +export type LpPositionsByWalletQueryType = z.infer< + typeof LpPositionsByWalletQuerySchema +>; diff --git a/src/modules/lp/lp.service.ts b/src/modules/lp/lp.service.ts new file mode 100644 index 00000000..e1a8103b --- /dev/null +++ b/src/modules/lp/lp.service.ts @@ -0,0 +1,114 @@ +// src/modules/lp/lp.service.ts +// +// LP position and rewards API (#980). + +import { prisma } from '../../utils/prisma.utils'; + +export class LpPositionNotFoundError extends Error { + constructor(lpId: string) { + super(`LP position not found: ${lpId}`); + this.name = 'LpPositionNotFoundError'; + } +} + +export interface LpPositionSummary { + lpId: string; + wallet: string; + keyId: string; + sharePercent: string; + accruedRewards: string; + status: string; + createdAt: string; + updatedAt: string; +} + +function toSummary(position: { + lpId: string; + wallet: string; + keyId: string; + sharePercent: unknown; + accruedRewards: unknown; + status: string; + createdAt: Date; + updatedAt: Date; +}): LpPositionSummary { + return { + lpId: position.lpId, + wallet: position.wallet, + keyId: position.keyId, + sharePercent: String(position.sharePercent), + accruedRewards: String(position.accruedRewards), + status: position.status, + createdAt: position.createdAt.toISOString(), + updatedAt: position.updatedAt.toISOString(), + }; +} + +/** All active LP positions for a wallet. */ +export async function getLpPositionsByWallet( + wallet: string +): Promise { + const positions = await prisma.lpPosition.findMany({ + where: { wallet, status: 'active' }, + orderBy: { createdAt: 'desc' }, + }); + return positions.map(toSummary); +} + +/** + * Single LP position by lpId, scoped to the owning wallet. Throws + * LpPositionNotFoundError both when the position doesn't exist and when it + * exists but isn't owned by `wallet`, so ownership is never leaked via a + * distinct error type/status. + */ +export async function getLpPositionById( + lpId: string, + wallet: string +): Promise { + const position = await prisma.lpPosition.findUnique({ where: { lpId } }); + if (!position || position.wallet !== wallet) { + throw new LpPositionNotFoundError(lpId); + } + return toSummary(position); +} + +export interface LpPoolSummary { + keyId: string; + totalPoolSize: string; + /** + * Simplified APR estimate: recent accrued rewards across the pool, + * annualized against pool size. This is NOT a precise on-chain yield + * figure — it ignores compounding, time-weighting of individual + * positions, and reward-rate changes over time. Treat as a rough + * indicator only. + */ + estimatedApr: number; +} + +/** Total active pool size and a simplified APR estimate for a key. */ +export async function getLpPoolSummary(keyId: string): Promise { + const positions = await prisma.lpPosition.findMany({ + where: { keyId, status: 'active' }, + }); + + const totalPoolSize = positions.reduce( + (sum, p) => sum + Number(p.sharePercent), + 0 + ); + const totalRewards = positions.reduce( + (sum, p) => sum + Number(p.accruedRewards), + 0 + ); + + // Simplified placeholder: annualize accrued-to-date rewards against pool + // size as if the current accrual rate held for a full year. Real APR + // would need a time-windowed reward rate, not lifetime accrued rewards. + const estimatedApr = + totalPoolSize > 0 ? (totalRewards / totalPoolSize) * 100 : 0; + + return { + keyId, + totalPoolSize: String(totalPoolSize), + estimatedApr, + }; +}