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/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/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/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/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/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/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/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/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; + } + }); +} 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; + } + }); +} 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). 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, + }; +}