diff --git a/.env.example b/.env.example index fd8463b7..b6373ba7 100644 --- a/.env.example +++ b/.env.example @@ -52,6 +52,10 @@ REVENUE_DISTRIBUTION_CYCLE_DAYS=7 # Comma-separated Stellar addresses of the three admin wallets. ADMIN_MULTISIG_WALLETS= +# Key sunset watch (#931): days of inactivity before a key is "near threshold" +# in GET /keys/sunset-watch. Keys flagged on-chain are always included. +KEY_SUNSET_INACTIVITY_THRESHOLD_DAYS=30 + # Creator and indexer tuning INDEXER_JITTER_FACTOR=0.1 BACKGROUND_JOB_LOCK_TTL_MS=300000 diff --git a/prisma/schema/creator.prisma b/prisma/schema/creator.prisma index e5237ff8..8f24e81c 100644 --- a/prisma/schema/creator.prisma +++ b/prisma/schema/creator.prisma @@ -31,6 +31,9 @@ model CreatorProfile { buybackPriceXlm Decimal? @db.Decimal(20, 7) /// Buyback window expiry; buybacks are rejected with 410 after this time (#882). buybackExpiresAt DateTime? + /// When the on-chain KeySunsetFlagged event was processed for this key (#931). + /// Null while the key has not been flagged for sunset. + sunsetFlaggedAt DateTime? /// Buy cooldown between purchases by the same wallet, in ledgers (0 = no cooldown). cooldownLedgers Int @default(0) perks Json? diff --git a/prisma/schema/governance-delegation.prisma b/prisma/schema/governance-delegation.prisma new file mode 100644 index 00000000..46a2c96f --- /dev/null +++ b/prisma/schema/governance-delegation.prisma @@ -0,0 +1,80 @@ +// prisma/schema/governance-delegation.prisma +// Delegation read model for governance voting (#933). +// +// VoteDelegation is the current-state table: one row per (delegator, keyId) +// pair, tracking the active delegate and whether the delegation has been +// revoked. Only rows where isActive = true count toward vote weight. +// +// VoteDelegationHistory is the append-only audit log: every DelegationSet and +// DelegationRevoked on-chain event appends a row so +// GET /governance/delegation/:wallet/history can return full provenance. +// +// Idempotency keys: +// DelegationSet — unique(txHash, eventIndex) on VoteDelegationHistory +// DelegationRevoked — same unique(txHash, eventIndex) on VoteDelegationHistory; +// the SET handler also guards on VoteDelegation.txHash +// for the upsert. + +model VoteDelegation { + id String @id @default(cuid()) + + /// Wallet that is delegating its vote weight. + delegatorWallet String + + /// Creator key scope this delegation applies to. + /// A wallet may hold different keys and delegate each independently. + keyId String + + /// Wallet receiving the delegated weight. + delegateeWallet String + + /// False once a DelegationRevoked event has been processed. + isActive Boolean @default(true) + + /// Ledger of the most recent state-changing event (set or revoke). + ledger Int + + /// Transaction hash of the most recent state-changing event. + /// Used as an optimistic concurrency guard so out-of-order replays + /// cannot clobber a newer state with an older one. + txHash String + + /// On-chain timestamp of the most recent state-changing event. + occurredAt DateTime + + createdAt DateTime @default(now()) + updatedAt DateTime @updatedAt + + @@unique([delegatorWallet, keyId]) + @@index([delegateeWallet, isActive]) + @@index([delegatorWallet]) + @@index([keyId, isActive]) + @@map("vote_delegations") +} + +/// Append-only log of every delegation change event. +/// One row per DelegationSet or DelegationRevoked contract event. +model VoteDelegationHistory { + id String @id @default(cuid()) + + delegatorWallet String + keyId String + + /// For DelegationSet events; null for DelegationRevoked entries. + delegateeWallet String? + + /// 'set' | 'revoked' + action String + + ledger Int + txHash String + eventIndex Int + occurredAt DateTime + createdAt DateTime @default(now()) + + @@unique([txHash, eventIndex]) + @@index([delegatorWallet, occurredAt(sort: Desc)]) + @@index([delegateeWallet, occurredAt(sort: Desc)]) + @@index([keyId]) + @@map("vote_delegation_history") +} diff --git a/prisma/schema/migrations/20260925000100_add_key_sunset_flagged_at/migration.sql b/prisma/schema/migrations/20260925000100_add_key_sunset_flagged_at/migration.sql new file mode 100644 index 00000000..01f4aaf6 --- /dev/null +++ b/prisma/schema/migrations/20260925000100_add_key_sunset_flagged_at/migration.sql @@ -0,0 +1,12 @@ +-- Migration: add_key_sunset_flagged_at (#931) +-- +-- Adds sunsetFlaggedAt to CreatorProfile so the indexer can persist the +-- timestamp when a KeySunsetFlagged on-chain event is processed. The column +-- is nullable; NULL means the key has never been flagged for sunset. +-- +-- Also extends the ActivityType enum with KEY_SUNSET_FLAGGED so the event +-- can be written to the Activity audit trail. + +ALTER TABLE "CreatorProfile" ADD COLUMN "sunsetFlaggedAt" TIMESTAMP(3); + +ALTER TYPE "ActivityType" ADD VALUE IF NOT EXISTS 'KEY_SUNSET_FLAGGED'; diff --git a/prisma/schema/migrations/20260925000200_add_staking_nfts/migration.sql b/prisma/schema/migrations/20260925000200_add_staking_nfts/migration.sql new file mode 100644 index 00000000..9460e898 --- /dev/null +++ b/prisma/schema/migrations/20260925000200_add_staking_nfts/migration.sql @@ -0,0 +1,69 @@ +-- Migration: add_staking_nfts (#932) +-- +-- Creates the staking NFT read model: +-- staking_nfts — one row per minted stake receipt NFT +-- staking_nft_transfers — append-only ownership history per NFT +-- +-- Also extends the ActivityType enum with the two new staking event types so +-- the indexer can write audit-trail Activity records. + +-- staking_nfts ────────────────────────────────────────────────────────────── + +CREATE TABLE "staking_nfts" ( + "id" TEXT NOT NULL, + "tokenId" TEXT NOT NULL, + "ownerAddress" TEXT NOT NULL, + "keyId" TEXT NOT NULL, + "stakedAmount" DECIMAL(30, 7) NOT NULL, + "lockExpiryLedger" INTEGER, + "lockExpiresAt" TIMESTAMP(3), + "burned" BOOLEAN NOT NULL DEFAULT FALSE, + "burnedAt" TIMESTAMP(3), + "mintLedger" INTEGER NOT NULL, + "mintTxHash" TEXT NOT NULL, + "createdAt" TIMESTAMP(3) NOT NULL DEFAULT CURRENT_TIMESTAMP, + "updatedAt" TIMESTAMP(3) NOT NULL, + + CONSTRAINT "staking_nfts_pkey" PRIMARY KEY ("id") +); + +CREATE UNIQUE INDEX "staking_nfts_tokenId_key" ON "staking_nfts" ("tokenId"); +CREATE UNIQUE INDEX "staking_nfts_mintTxHash_key" ON "staking_nfts" ("mintTxHash"); +CREATE INDEX "staking_nfts_ownerAddress_idx" ON "staking_nfts" ("ownerAddress"); +CREATE INDEX "staking_nfts_keyId_idx" ON "staking_nfts" ("keyId"); +CREATE INDEX "staking_nfts_ownerAddress_burned_idx" ON "staking_nfts" ("ownerAddress", "burned"); + +-- staking_nft_transfers ───────────────────────────────────────────────────── + +CREATE TABLE "staking_nft_transfers" ( + "id" TEXT NOT NULL, + "nftId" TEXT NOT NULL, + "fromAddress" TEXT, + "toAddress" TEXT NOT NULL, + "ledger" INTEGER NOT NULL, + "txHash" TEXT NOT NULL, + "eventIndex" INTEGER NOT NULL, + "occurredAt" TIMESTAMP(3) NOT NULL, + "createdAt" TIMESTAMP(3) NOT NULL DEFAULT CURRENT_TIMESTAMP, + + CONSTRAINT "staking_nft_transfers_pkey" PRIMARY KEY ("id") +); + +CREATE UNIQUE INDEX "staking_nft_transfers_txHash_eventIndex_key" + ON "staking_nft_transfers" ("txHash", "eventIndex"); + +CREATE INDEX "staking_nft_transfers_nftId_occurredAt_idx" + ON "staking_nft_transfers" ("nftId", "occurredAt" DESC); + +CREATE INDEX "staking_nft_transfers_toAddress_idx" + ON "staking_nft_transfers" ("toAddress"); + +ALTER TABLE "staking_nft_transfers" + ADD CONSTRAINT "staking_nft_transfers_nftId_fkey" + FOREIGN KEY ("nftId") REFERENCES "staking_nfts" ("id") + ON DELETE CASCADE ON UPDATE CASCADE; + +-- ActivityType enum extensions ────────────────────────────────────────────── + +ALTER TYPE "ActivityType" ADD VALUE IF NOT EXISTS 'STAKE_NFT_MINTED'; +ALTER TYPE "ActivityType" ADD VALUE IF NOT EXISTS 'STAKE_NFT_TRANSFERRED'; diff --git a/prisma/schema/migrations/20260925000300_add_governance_delegation/migration.sql b/prisma/schema/migrations/20260925000300_add_governance_delegation/migration.sql new file mode 100644 index 00000000..6d0453cb --- /dev/null +++ b/prisma/schema/migrations/20260925000300_add_governance_delegation/migration.sql @@ -0,0 +1,71 @@ +-- Migration: add_governance_delegation (#933) +-- +-- Creates: +-- vote_delegations — current-state table (one row per delegator/keyId pair) +-- vote_delegation_history — append-only audit log of every delegation event +-- +-- Extends ActivityType enum with the three new delegation/vote event values. + +-- vote_delegations ────────────────────────────────────────────────────────── + +CREATE TABLE "vote_delegations" ( + "id" TEXT NOT NULL, + "delegatorWallet" TEXT NOT NULL, + "keyId" TEXT NOT NULL, + "delegateeWallet" TEXT NOT NULL, + "isActive" BOOLEAN NOT NULL DEFAULT TRUE, + "ledger" INTEGER NOT NULL, + "txHash" TEXT NOT NULL, + "occurredAt" TIMESTAMP(3) NOT NULL, + "createdAt" TIMESTAMP(3) NOT NULL DEFAULT CURRENT_TIMESTAMP, + "updatedAt" TIMESTAMP(3) NOT NULL, + + CONSTRAINT "vote_delegations_pkey" PRIMARY KEY ("id") +); + +CREATE UNIQUE INDEX "vote_delegations_delegatorWallet_keyId_key" + ON "vote_delegations" ("delegatorWallet", "keyId"); + +CREATE INDEX "vote_delegations_delegateeWallet_isActive_idx" + ON "vote_delegations" ("delegateeWallet", "isActive"); + +CREATE INDEX "vote_delegations_delegatorWallet_idx" + ON "vote_delegations" ("delegatorWallet"); + +CREATE INDEX "vote_delegations_keyId_isActive_idx" + ON "vote_delegations" ("keyId", "isActive"); + +-- vote_delegation_history ─────────────────────────────────────────────────── + +CREATE TABLE "vote_delegation_history" ( + "id" TEXT NOT NULL, + "delegatorWallet" TEXT NOT NULL, + "keyId" TEXT NOT NULL, + "delegateeWallet" TEXT, + "action" TEXT NOT NULL, + "ledger" INTEGER NOT NULL, + "txHash" TEXT NOT NULL, + "eventIndex" INTEGER NOT NULL, + "occurredAt" TIMESTAMP(3) NOT NULL, + "createdAt" TIMESTAMP(3) NOT NULL DEFAULT CURRENT_TIMESTAMP, + + CONSTRAINT "vote_delegation_history_pkey" PRIMARY KEY ("id") +); + +CREATE UNIQUE INDEX "vote_delegation_history_txHash_eventIndex_key" + ON "vote_delegation_history" ("txHash", "eventIndex"); + +CREATE INDEX "vote_delegation_history_delegatorWallet_occurredAt_idx" + ON "vote_delegation_history" ("delegatorWallet", "occurredAt" DESC); + +CREATE INDEX "vote_delegation_history_delegateeWallet_occurredAt_idx" + ON "vote_delegation_history" ("delegateeWallet", "occurredAt" DESC); + +CREATE INDEX "vote_delegation_history_keyId_idx" + ON "vote_delegation_history" ("keyId"); + +-- ActivityType enum extensions ────────────────────────────────────────────── + +ALTER TYPE "ActivityType" ADD VALUE IF NOT EXISTS 'GOVERNANCE_VOTE_CAST'; +ALTER TYPE "ActivityType" ADD VALUE IF NOT EXISTS 'GOVERNANCE_DELEGATION_SET'; +ALTER TYPE "ActivityType" ADD VALUE IF NOT EXISTS 'GOVERNANCE_DELEGATION_REVOKED'; diff --git a/prisma/schema/migrations/20260925000400_add_trade_payment_asset/migration.sql b/prisma/schema/migrations/20260925000400_add_trade_payment_asset/migration.sql new file mode 100644 index 00000000..1fa71196 --- /dev/null +++ b/prisma/schema/migrations/20260925000400_add_trade_payment_asset/migration.sql @@ -0,0 +1,11 @@ +-- Migration: add_trade_payment_asset (#934) +-- +-- Adds paymentAsset to the Trade table to track which asset was used +-- for each key purchase. Existing rows default to 'XLM'. +-- Adds targeted indexes for aggregation queries. + +ALTER TABLE "Trade" + ADD COLUMN "paymentAsset" TEXT NOT NULL DEFAULT 'XLM'; + +CREATE INDEX "Trade_paymentAsset_idx" ON "Trade" ("paymentAsset"); +CREATE INDEX "Trade_creatorId_paymentAsset_idx" ON "Trade" ("creatorId", "paymentAsset"); diff --git a/prisma/schema/staking.prisma b/prisma/schema/staking.prisma index 241d6714..9aed977c 100644 --- a/prisma/schema/staking.prisma +++ b/prisma/schema/staking.prisma @@ -1,5 +1,6 @@ // prisma/schema/staking.prisma // Staking reward multiplier tiers and positions configuration (#942). +// StakingNft and StakingNftTransfer models (#932) are also defined here. model StakingMultiplierTier { id String @id @default(cuid()) @@ -31,3 +32,53 @@ model StakingPosition { @@index([keyId]) @@map("staking_positions") } + +/// Tradeable stake receipt NFT (#932). +model StakingNft { + id String @id @default(cuid()) + /// On-chain token ID (unique per network). + tokenId String @unique + /// Current holder wallet. + ownerAddress String + /// Creator key the staked tokens belong to. + keyId String + /// Number of creator keys locked in this position. + stakedAmount Decimal @db.Decimal(30, 7) + lockExpiryLedger Int? + lockExpiresAt DateTime? + burned Boolean @default(false) + burnedAt DateTime? + mintLedger Int + /// Mint tx hash — idempotency key for the indexer. + mintTxHash String @unique + createdAt DateTime @default(now()) + updatedAt DateTime @updatedAt + + transfers StakingNftTransfer[] + + @@index([ownerAddress]) + @@index([keyId]) + @@index([ownerAddress, burned]) + @@map("staking_nfts") +} + +/// Append-only ownership history for a staking NFT (#932). +model StakingNftTransfer { + id String @id @default(cuid()) + nftId String + /// Null on the initial mint. + fromAddress String? + toAddress String + ledger Int + txHash String + eventIndex Int + occurredAt DateTime + createdAt DateTime @default(now()) + + nft StakingNft @relation(fields: [nftId], references: [id], onDelete: Cascade) + + @@unique([txHash, eventIndex]) + @@index([nftId, occurredAt(sort: Desc)]) + @@index([toAddress]) + @@map("staking_nft_transfers") +} diff --git a/prisma/schema/trade.prisma b/prisma/schema/trade.prisma index 6883cb94..d59ee765 100644 --- a/prisma/schema/trade.prisma +++ b/prisma/schema/trade.prisma @@ -1,18 +1,23 @@ // prisma/schema/trade.prisma model Trade { - id String @id @default(cuid()) - buyer String - creatorId String - quantity String - price String - ledger Int - txHash String - timestamp DateTime - createdAt DateTime @default(now()) + id String @id @default(cuid()) + buyer String + creatorId String + quantity String + price String + ledger Int + txHash String + timestamp DateTime + /// Payment asset used for this purchase (e.g. 'XLM', 'USDC'). + /// Defaults to 'XLM' for trades recorded before multi-currency support (#934). + paymentAsset String @default("XLM") + createdAt DateTime @default(now()) @@unique([ledger, txHash]) @@index([creatorId]) @@index([buyer]) @@index([ledger]) + @@index([paymentAsset]) + @@index([creatorId, paymentAsset]) } diff --git a/src/config.schema.ts b/src/config.schema.ts index aff00a88..1f09a4f1 100644 --- a/src/config.schema.ts +++ b/src/config.schema.ts @@ -278,6 +278,16 @@ export const envSchema = z .positive() .default(5), + // Key sunset watch (#931): number of consecutive inactive days before a + // key is considered "near threshold" and surfaced by GET /keys/sunset-watch. + // Keys whose on-chain KeySunsetFlagged event has been processed always + // appear regardless of this threshold. Defaults to 30 days. + KEY_SUNSET_INACTIVITY_THRESHOLD_DAYS: z.coerce + .number() + .int() + .positive() + .default(30), + // Request body size limits (see docs/body-size-limits.md). // Accepts any size string understood by the `bytes` package used // internally by body-parser (e.g. '100kb', '1mb', '10mb'). diff --git a/src/modules/governance/governance-delegation-indexer.service.ts b/src/modules/governance/governance-delegation-indexer.service.ts new file mode 100644 index 00000000..472dfca7 --- /dev/null +++ b/src/modules/governance/governance-delegation-indexer.service.ts @@ -0,0 +1,438 @@ +// src/modules/governance/governance-delegation-indexer.service.ts +// Indexer event handlers for governance delegation contract events (#933). +// +// Two event types are handled: +// +// DELEGATION_SET +// Emitted when a wallet sets (or updates) its vote delegate for a key. +// Upserts VoteDelegation with isActive=true and appends a +// VoteDelegationHistory row with action='set'. +// Guard: if the stored ledger is higher than the incoming event, the +// event is from a replay of an older batch — skip it to avoid clobbering +// a newer revoke with an older set. +// +// DELEGATION_REVOKED +// Emitted when a wallet removes its active vote delegation. +// Sets isActive=false on the VoteDelegation row and appends a +// VoteDelegationHistory row with action='revoked'. +// Skip gracefully if no active delegation exists (out-of-order delivery). +// +// Both handlers: +// - Validate required fields and skip with a warn on missing data +// - Are idempotent via unique(txHash, eventIndex) on VoteDelegationHistory +// - Invalidate the delegation caches so reads reflect the change immediately +// - Write an Activity audit record + +import { prisma } from '../../utils/prisma.utils'; +import { logger } from '../../utils/logger.utils'; +import { + processIndexerChainEvents, + IndexerChainEvent, +} from '../../utils/indexer-event-processor.utils'; +import { cacheInvalidate } from '../../utils/redis.utils'; +import { + delegateCachePattern, + delegatorsCachePattern, + delegationHistoryCachePattern, +} from './governance-delegation.service'; + +// ── Typed event interfaces ──────────────────────────────────── + +/** + * Contract event emitted when a wallet sets or updates its vote delegate. + * + * Required fields: + * delegatorWallet — wallet that is delegating + * delegateeWallet — wallet receiving the delegation + * keyId — creator key scope + * ledger — ledger sequence + * txHash — transaction hash + * eventIndex — position within the transaction (dedup) + * occurredAt — ISO-8601 timestamp + */ +export interface DelegationSetEvent extends IndexerChainEvent { + eventType: 'DELEGATION_SET'; + delegatorWallet: string; + delegateeWallet: string; + keyId: string; + occurredAt: string; +} + +/** + * Contract event emitted when a wallet revokes its active vote delegation. + * + * Required fields: + * delegatorWallet — wallet revoking the delegation + * keyId — creator key scope + * ledger — ledger sequence + * txHash — transaction hash + * eventIndex — position within the transaction (dedup) + * occurredAt — ISO-8601 timestamp + */ +export interface DelegationRevokedEvent extends IndexerChainEvent { + eventType: 'DELEGATION_REVOKED'; + delegatorWallet: string; + keyId: string; + occurredAt: string; +} + +// ── DELEGATION_SET ──────────────────────────────────────────── + +const DELEGATION_SET_REQUIRED_FIELDS: (keyof DelegationSetEvent)[] = [ + 'delegatorWallet', + 'delegateeWallet', + 'keyId', + 'occurredAt', + 'ledger', + 'txHash', + 'eventIndex', +]; + +/** + * Process a batch of DELEGATION_SET events. + * + * For each event: + * 1. Validates required fields. + * 2. Idempotency: skips if VoteDelegationHistory already has a row for + * this (txHash, eventIndex). + * 3. Ledger guard: skips if the stored VoteDelegation.ledger is higher + * (event is a stale replay). + * 4. Upserts VoteDelegation (delegateeWallet, isActive=true, ledger, txHash). + * 5. Appends a VoteDelegationHistory row (action='set'). + * 6. Writes an Activity record. + * 7. Invalidates caches for both wallets. + */ +export async function processDelegationSetEvents( + events: IndexerChainEvent[] +): Promise { + await processIndexerChainEvents(events, async event => { + if (event.eventType !== 'DELEGATION_SET') return; + + const e = event as DelegationSetEvent; + + for (const field of DELEGATION_SET_REQUIRED_FIELDS) { + const value = e[field]; + if (value === undefined || value === null || value === '') { + logger.warn( + { + eventId: `${e.txHash}:${e.eventIndex}`, + missingField: field, + }, + 'Skipping DELEGATION_SET event due to missing required field' + ); + return; + } + } + + // Idempotency: skip if this exact event is already recorded. + const existingHistory = await prisma.voteDelegationHistory.findUnique({ + where: { + txHash_eventIndex: { + txHash: e.txHash, + eventIndex: Number(e.eventIndex), + }, + }, + select: { id: true }, + }); + if (existingHistory) { + logger.info( + { eventId: `${e.txHash}:${e.eventIndex}`, delegatorWallet: e.delegatorWallet }, + 'DELEGATION_SET already recorded; skipping duplicate event' + ); + return; + } + + const occurredAt = new Date(e.occurredAt); + const incomingLedger = Number(e.ledger); + + // Ledger guard: if the current state was set by a later ledger, this + // is a stale replay of an earlier event — skip to avoid regression. + const existing = await prisma.voteDelegation.findUnique({ + where: { + delegatorWallet_keyId: { + delegatorWallet: e.delegatorWallet, + keyId: e.keyId, + }, + }, + select: { ledger: true, delegateeWallet: true }, + }); + + if (existing && existing.ledger > incomingLedger) { + logger.info( + { + eventId: `${e.txHash}:${e.eventIndex}`, + delegatorWallet: e.delegatorWallet, + storedLedger: existing.ledger, + incomingLedger, + }, + 'DELEGATION_SET skipped: stored state is from a later ledger' + ); + // Still append history so the audit trail is complete. + await prisma.voteDelegationHistory.create({ + data: { + delegatorWallet: e.delegatorWallet, + keyId: e.keyId, + delegateeWallet: e.delegateeWallet, + action: 'set', + ledger: incomingLedger, + txHash: e.txHash, + eventIndex: Number(e.eventIndex), + occurredAt, + }, + }); + return; + } + + const previousDelegatee = existing?.delegateeWallet ?? null; + + await prisma.$transaction([ + // Upsert the current-state delegation row. + prisma.voteDelegation.upsert({ + where: { + delegatorWallet_keyId: { + delegatorWallet: e.delegatorWallet, + keyId: e.keyId, + }, + }, + update: { + delegateeWallet: e.delegateeWallet, + isActive: true, + ledger: incomingLedger, + txHash: e.txHash, + occurredAt, + }, + create: { + delegatorWallet: e.delegatorWallet, + keyId: e.keyId, + delegateeWallet: e.delegateeWallet, + isActive: true, + ledger: incomingLedger, + txHash: e.txHash, + occurredAt, + }, + }), + // Append history row. + prisma.voteDelegationHistory.create({ + data: { + delegatorWallet: e.delegatorWallet, + keyId: e.keyId, + delegateeWallet: e.delegateeWallet, + action: 'set', + ledger: incomingLedger, + txHash: e.txHash, + eventIndex: Number(e.eventIndex), + occurredAt, + }, + }), + // Audit trail activity record. + prisma.activity.create({ + data: { + type: 'GOVERNANCE_DELEGATION_SET' as any, + actor: e.delegatorWallet, + target: e.delegateeWallet, + creatorId: e.keyId, + payload: { + delegatorWallet: e.delegatorWallet, + delegateeWallet: e.delegateeWallet, + previousDelegatee, + keyId: e.keyId, + ledger_sequence: incomingLedger, + }, + createdAt: occurredAt, + }, + }), + ]); + + // Invalidate caches for both wallets (delegator's "who I delegated to" + // and delegatee's "who delegated to me", plus any previous delegatee). + const patternsToInvalidate = [ + delegateCachePattern(e.delegatorWallet), + delegatorsCachePattern(e.delegateeWallet), + delegationHistoryCachePattern(e.delegatorWallet), + delegationHistoryCachePattern(e.delegateeWallet), + ]; + if (previousDelegatee && previousDelegatee !== e.delegateeWallet) { + patternsToInvalidate.push(delegatorsCachePattern(previousDelegatee)); + patternsToInvalidate.push(delegationHistoryCachePattern(previousDelegatee)); + } + await cacheInvalidate(...patternsToInvalidate); + + logger.info( + { + delegatorWallet: e.delegatorWallet, + delegateeWallet: e.delegateeWallet, + keyId: e.keyId, + ledger: incomingLedger, + txHash: e.txHash, + }, + 'DELEGATION_SET event processed' + ); + }); +} + +// ── DELEGATION_REVOKED ──────────────────────────────────────── + +const DELEGATION_REVOKED_REQUIRED_FIELDS: (keyof DelegationRevokedEvent)[] = [ + 'delegatorWallet', + 'keyId', + 'occurredAt', + 'ledger', + 'txHash', + 'eventIndex', +]; + +/** + * Process a batch of DELEGATION_REVOKED events. + * + * For each event: + * 1. Validates required fields. + * 2. Idempotency: skips if this (txHash, eventIndex) is already recorded. + * 3. Resolves the current delegation row; skips if none exists or already + * revoked (out-of-order delivery handled gracefully). + * 4. Sets VoteDelegation.isActive=false. + * 5. Appends a VoteDelegationHistory row (action='revoked'). + * 6. Writes an Activity record. + * 7. Invalidates caches for both wallets. + */ +export async function processDelegationRevokedEvents( + events: IndexerChainEvent[] +): Promise { + await processIndexerChainEvents(events, async event => { + if (event.eventType !== 'DELEGATION_REVOKED') return; + + const e = event as DelegationRevokedEvent; + + for (const field of DELEGATION_REVOKED_REQUIRED_FIELDS) { + const value = e[field]; + if (value === undefined || value === null || value === '') { + logger.warn( + { + eventId: `${e.txHash}:${e.eventIndex}`, + missingField: field, + }, + 'Skipping DELEGATION_REVOKED event due to missing required field' + ); + return; + } + } + + // Idempotency check. + const existingHistory = await prisma.voteDelegationHistory.findUnique({ + where: { + txHash_eventIndex: { + txHash: e.txHash, + eventIndex: Number(e.eventIndex), + }, + }, + select: { id: true }, + }); + if (existingHistory) { + logger.info( + { eventId: `${e.txHash}:${e.eventIndex}`, delegatorWallet: e.delegatorWallet }, + 'DELEGATION_REVOKED already recorded; skipping duplicate event' + ); + return; + } + + // Resolve the current delegation. + const delegation = await prisma.voteDelegation.findUnique({ + where: { + delegatorWallet_keyId: { + delegatorWallet: e.delegatorWallet, + keyId: e.keyId, + }, + }, + select: { id: true, delegateeWallet: true, isActive: true }, + }); + + const occurredAt = new Date(e.occurredAt); + const incomingLedger = Number(e.ledger); + + if (!delegation) { + // No delegation exists — still record history for auditability. + logger.warn( + { + eventId: `${e.txHash}:${e.eventIndex}`, + delegatorWallet: e.delegatorWallet, + keyId: e.keyId, + }, + 'DELEGATION_REVOKED references unknown delegation; recording history only' + ); + await prisma.voteDelegationHistory.create({ + data: { + delegatorWallet: e.delegatorWallet, + keyId: e.keyId, + delegateeWallet: null, + action: 'revoked', + ledger: incomingLedger, + txHash: e.txHash, + eventIndex: Number(e.eventIndex), + occurredAt, + }, + }); + return; + } + + const previousDelegatee = delegation.delegateeWallet; + + await prisma.$transaction([ + // Mark the delegation inactive. + prisma.voteDelegation.update({ + where: { id: delegation.id }, + data: { + isActive: false, + ledger: incomingLedger, + txHash: e.txHash, + occurredAt, + }, + }), + // Append history row. + prisma.voteDelegationHistory.create({ + data: { + delegatorWallet: e.delegatorWallet, + keyId: e.keyId, + delegateeWallet: null, + action: 'revoked', + ledger: incomingLedger, + txHash: e.txHash, + eventIndex: Number(e.eventIndex), + occurredAt, + }, + }), + // Audit trail. + prisma.activity.create({ + data: { + type: 'GOVERNANCE_DELEGATION_REVOKED' as any, + actor: e.delegatorWallet, + target: previousDelegatee, + creatorId: e.keyId, + payload: { + delegatorWallet: e.delegatorWallet, + revokedDelegatee: previousDelegatee, + keyId: e.keyId, + ledger_sequence: incomingLedger, + }, + createdAt: occurredAt, + }, + }), + ]); + + await cacheInvalidate( + delegateCachePattern(e.delegatorWallet), + delegatorsCachePattern(previousDelegatee), + delegationHistoryCachePattern(e.delegatorWallet), + delegationHistoryCachePattern(previousDelegatee) + ); + + logger.info( + { + delegatorWallet: e.delegatorWallet, + revokedDelegatee: previousDelegatee, + keyId: e.keyId, + ledger: incomingLedger, + txHash: e.txHash, + }, + 'DELEGATION_REVOKED event processed' + ); + }); +} diff --git a/src/modules/governance/governance-delegation.service.ts b/src/modules/governance/governance-delegation.service.ts new file mode 100644 index 00000000..e08646e3 --- /dev/null +++ b/src/modules/governance/governance-delegation.service.ts @@ -0,0 +1,347 @@ +// src/modules/governance/governance-delegation.service.ts +// Read-model service for delegated governance voting (#933). +// +// All writes are performed by the indexer +// (governance-delegation-indexer.service.ts). This module is read-only. +// +// Three query surfaces: +// getCurrentDelegate — who a wallet has delegated to (per key scope) +// getActiveDelegators — wallets that have delegated TO a given address +// getDelegationHistory — full event log for a wallet across all keys + +import { prisma } from '../../utils/prisma.utils'; +import { cacheGetJson, cacheSetJson } from '../../utils/redis.utils'; +import { buildOffsetPaginationMeta } from '../../utils/pagination.utils'; + +// ── Cache TTLs ──────────────────────────────────────────────── + +const DELEGATE_CACHE_TTL_SECONDS = 60; +const DELEGATORS_CACHE_TTL_SECONDS = 60; +const HISTORY_CACHE_TTL_SECONDS = 60; + +// ── Error classes ───────────────────────────────────────────── + +export class DelegationNotFoundError extends Error { + constructor(wallet: string, keyId?: string) { + super( + keyId + ? `No active delegation found for wallet ${wallet} on key ${keyId}` + : `No active delegation found for wallet ${wallet}` + ); + this.name = 'DelegationNotFoundError'; + } +} + +// ── Shared item shapes ──────────────────────────────────────── + +export interface DelegationItem { + id: string; + delegatorWallet: string; + keyId: string; + delegateeWallet: string; + isActive: boolean; + ledger: number; + occurredAt: string; + createdAt: string; + updatedAt: string; +} + +export interface DelegatorItem { + id: string; + delegatorWallet: string; + keyId: string; + ledger: number; + occurredAt: string; +} + +export interface DelegationHistoryItem { + id: string; + delegatorWallet: string; + keyId: string; + /** null for 'revoked' entries */ + delegateeWallet: string | null; + action: 'set' | 'revoked'; + ledger: number; + txHash: string; + occurredAt: string; +} + +// ── Cache invalidation helpers (used by indexer) ────────────── + +export function delegateCachePattern(wallet: string): string { + return `governance:delegation:delegate:${wallet}:*`; +} + +export function delegatorsCachePattern(wallet: string): string { + return `governance:delegation:delegators:${wallet}:*`; +} + +export function delegationHistoryCachePattern(wallet: string): string { + return `governance:delegation:history:${wallet}:*`; +} + +// ── Current delegate ────────────────────────────────────────── + +export interface GetCurrentDelegateQuery { + /** Delegator wallet to look up. */ + wallet: string; + /** Optional: narrow to a specific creator key. */ + keyId?: string; +} + +export type CurrentDelegateResult = + | { delegated: false } + | { delegated: true; delegation: DelegationItem }; + +/** + * Return the current active delegate for a wallet, optionally scoped to a + * specific creator key. + * + * When `keyId` is omitted and the wallet has multiple active delegations + * across different keys, the most-recently-set one is returned. Callers + * that need per-key precision should always supply `keyId`. + * + * Returns `{ delegated: false }` when no active delegation exists rather than + * throwing, so the route can return a clean 200 with that shape. + * + * Cached per (wallet, keyId) for 60 s. + */ +export async function getCurrentDelegate( + query: GetCurrentDelegateQuery +): Promise { + const { wallet, keyId } = query; + const cacheKey = `governance:delegation:delegate:${wallet}:${keyId ?? '_all'}`; + const cached = await cacheGetJson(cacheKey); + if (cached) return cached; + + const row = await prisma.voteDelegation.findFirst({ + where: { + delegatorWallet: wallet, + isActive: true, + ...(keyId ? { keyId } : {}), + }, + orderBy: { occurredAt: 'desc' }, + }); + + const result: CurrentDelegateResult = row + ? { delegated: true, delegation: mapDelegationRow(row) } + : { delegated: false }; + + await cacheSetJson(cacheKey, result, DELEGATE_CACHE_TTL_SECONDS); + return result; +} + +// ── Active delegators list ──────────────────────────────────── + +export interface GetActiveDelegatorsQuery { + /** Delegatee wallet: return wallets that delegate TO this address. */ + wallet: string; + /** Optional: narrow to a specific creator key. */ + keyId?: string; + limit: number; + offset: number; +} + +export interface ActiveDelegatorsResult { + items: DelegatorItem[]; + meta: ReturnType; +} + +/** + * Return all wallets that currently have an active delegation pointing to + * `wallet`, sorted by delegation date descending. + * + * Optionally filtered by `keyId` to show delegators on a single key. + * + * Cached per (wallet, keyId, limit, offset) for 60 s. + */ +export async function getActiveDelegators( + query: GetActiveDelegatorsQuery +): Promise { + const { wallet, keyId, limit, offset } = query; + const cacheKey = `governance:delegation:delegators:${wallet}:${keyId ?? '_all'}:${limit}:${offset}`; + const cached = await cacheGetJson(cacheKey); + if (cached) return cached; + + const where = { + delegateeWallet: wallet, + isActive: true, + ...(keyId ? { keyId } : {}), + }; + + const [rows, total] = await Promise.all([ + prisma.voteDelegation.findMany({ + where, + orderBy: { occurredAt: 'desc' }, + skip: offset, + take: limit, + select: { + id: true, + delegatorWallet: true, + keyId: true, + ledger: true, + occurredAt: true, + }, + }), + prisma.voteDelegation.count({ where }), + ]); + + const result: ActiveDelegatorsResult = { + items: rows.map(r => ({ + id: r.id, + delegatorWallet: r.delegatorWallet, + keyId: r.keyId, + ledger: r.ledger, + occurredAt: r.occurredAt.toISOString(), + })), + meta: buildOffsetPaginationMeta({ limit, offset, total }), + }; + + await cacheSetJson(cacheKey, result, DELEGATORS_CACHE_TTL_SECONDS); + return result; +} + +// ── Delegation history ──────────────────────────────────────── + +export interface GetDelegationHistoryQuery { + /** Wallet whose delegation history to return (as delegator OR delegatee). */ + wallet: string; + /** Optional: narrow to a specific creator key. */ + keyId?: string; + limit: number; + offset: number; +} + +export interface DelegationHistoryResult { + items: DelegationHistoryItem[]; + meta: ReturnType; +} + +/** + * Return the full delegation event log for `wallet` (both sides: events where + * the wallet was the delegator or the delegatee), newest first. + * + * Optionally filtered by `keyId`. + * + * Cached per (wallet, keyId, limit, offset) for 60 s. + */ +export async function getDelegationHistory( + query: GetDelegationHistoryQuery +): Promise { + const { wallet, keyId, limit, offset } = query; + const cacheKey = `governance:delegation:history:${wallet}:${keyId ?? '_all'}:${limit}:${offset}`; + const cached = await cacheGetJson(cacheKey); + if (cached) return cached; + + const where = { + OR: [ + { delegatorWallet: wallet }, + { delegateeWallet: wallet }, + ], + ...(keyId ? { keyId } : {}), + }; + + const [rows, total] = await Promise.all([ + prisma.voteDelegationHistory.findMany({ + where, + orderBy: { occurredAt: 'desc' }, + skip: offset, + take: limit, + }), + prisma.voteDelegationHistory.count({ where }), + ]); + + const result: DelegationHistoryResult = { + items: rows.map(mapHistoryRow), + meta: buildOffsetPaginationMeta({ limit, offset, total }), + }; + + await cacheSetJson(cacheKey, result, HISTORY_CACHE_TTL_SECONDS); + return result; +} + +// ── Vote-weight helper (used by castKeyProposalVote) ────────── + +/** + * Return the sum of `KeyOwnership.balance` for every wallet that currently + * has an active delegation pointing to `delegateeWallet` on `keyId`. + * + * This is the delegated weight that the delegatee carries in addition to + * their own balance when casting a governance vote. + * + * Not cached — called inside a vote transaction where freshness is critical. + */ +export async function getDelegatedVoteWeight( + delegateeWallet: string, + keyId: string +): Promise { + // Fetch all active delegator wallets for this delegatee+key combination. + const activeDelegations = await prisma.voteDelegation.findMany({ + where: { delegateeWallet, keyId, isActive: true }, + select: { delegatorWallet: true }, + }); + + if (activeDelegations.length === 0) return 0; + + const delegatorAddresses = activeDelegations.map(d => d.delegatorWallet); + + // Sum the key balances of all delegators. + const agg = await prisma.keyOwnership.aggregate({ + where: { + ownerAddress: { in: delegatorAddresses }, + creatorId: keyId, + balance: { gt: 0 }, + }, + _sum: { balance: true }, + }); + + return Number(agg._sum.balance ?? 0); +} + +// ── Row mappers ─────────────────────────────────────────────── + +function mapDelegationRow(row: { + id: string; + delegatorWallet: string; + keyId: string; + delegateeWallet: string; + isActive: boolean; + ledger: number; + occurredAt: Date; + createdAt: Date; + updatedAt: Date; +}): DelegationItem { + return { + id: row.id, + delegatorWallet: row.delegatorWallet, + keyId: row.keyId, + delegateeWallet: row.delegateeWallet, + isActive: row.isActive, + ledger: row.ledger, + occurredAt: row.occurredAt.toISOString(), + createdAt: row.createdAt.toISOString(), + updatedAt: row.updatedAt.toISOString(), + }; +} + +function mapHistoryRow(row: { + id: string; + delegatorWallet: string; + keyId: string; + delegateeWallet: string | null; + action: string; + ledger: number; + txHash: string; + occurredAt: Date; +}): DelegationHistoryItem { + return { + id: row.id, + delegatorWallet: row.delegatorWallet, + keyId: row.keyId, + delegateeWallet: row.delegateeWallet, + action: row.action as 'set' | 'revoked', + ledger: row.ledger, + txHash: row.txHash, + occurredAt: row.occurredAt.toISOString(), + }; +} diff --git a/src/modules/indexer/indexer-pipeline.service.ts b/src/modules/indexer/indexer-pipeline.service.ts index b03ebe07..7557b722 100644 --- a/src/modules/indexer/indexer-pipeline.service.ts +++ b/src/modules/indexer/indexer-pipeline.service.ts @@ -77,6 +77,11 @@ export async function processTradeEvents( const { creatorId, actor, amount, price, feePaid, tradeAt, ledger } = event; + // payment_asset is optional — absent events default to 'XLM' (#934). + const paymentAsset: string = + typeof event.paymentAsset === 'string' && event.paymentAsset.trim() !== '' + ? event.paymentAsset.trim().toUpperCase() + : 'XLM'; // 1. Create corresponding Activity record await prisma.activity.create({ @@ -89,6 +94,7 @@ export async function processTradeEvents( price_at_trade: price.toString(), fee_paid: feePaid.toString(), ledger_sequence: Number(ledger), + payment_asset: paymentAsset, }, createdAt: new Date(tradeAt), }, @@ -212,3 +218,112 @@ function computeBatchHash( .digest('hex') .slice(0, 16); } + +/** + * Processes a batch of on-chain KeySunsetFlagged events (#931). + * + * Each event stamps `sunsetFlaggedAt` on the matching CreatorProfile and + * writes a KEY_SUNSET_FLAGGED Activity record so the flag appears in the + * audit trail. Already-flagged keys are skipped (idempotent). + * + * Expected event fields: + * - eventType : 'KEY_SUNSET_FLAGGED' + * - creatorId : creator profile ID the flag applies to + * - flaggedAt : ISO-8601 timestamp from the contract (optional; falls back to now) + * - ledger : ledger sequence number + * - txHash : transaction hash (for dedup) + * - eventIndex : position within the transaction (for dedup) + */ +export async function processSunsetFlaggedEvents( + events: IndexerChainEvent[] +): Promise { + await processIndexerChainEvents(events, async (event) => { + if (event.eventType !== 'KEY_SUNSET_FLAGGED') { + return; + } + + const requiredFields = ['creatorId', 'ledger']; + for (const field of requiredFields) { + if ( + event[field] === undefined || + event[field] === null || + event[field] === '' + ) { + logger.warn( + { + eventId: `${event.txHash}:${event.eventIndex}`, + missingField: field, + }, + 'Skipping KEY_SUNSET_FLAGGED event due to missing required field' + ); + return; + } + } + + const creatorId = String(event.creatorId); + const flaggedAt = event.flaggedAt + ? new Date(String(event.flaggedAt)) + : new Date(); + + // Resolve to the canonical profile ID (event may carry handle or id). + const profile = await prisma.creatorProfile.findFirst({ + where: { OR: [{ id: creatorId }, { handle: creatorId }] }, + select: { id: true, sunsetFlaggedAt: true }, + }); + + if (!profile) { + logger.warn( + { + eventId: `${event.txHash}:${event.eventIndex}`, + creatorId, + }, + 'KEY_SUNSET_FLAGGED event references unknown creator; skipping' + ); + return; + } + + // Idempotent: if the flag is already set, skip further writes. + if (profile.sunsetFlaggedAt !== null) { + logger.info( + { + eventId: `${event.txHash}:${event.eventIndex}`, + creatorId: profile.id, + sunsetFlaggedAt: profile.sunsetFlaggedAt.toISOString(), + }, + 'KEY_SUNSET_FLAGGED already recorded; skipping duplicate event' + ); + return; + } + + await prisma.$transaction([ + // Stamp the flag on the creator profile. + prisma.creatorProfile.update({ + where: { id: profile.id }, + data: { sunsetFlaggedAt: flaggedAt }, + }), + // Write an Activity record for the audit trail. + prisma.activity.create({ + data: { + type: 'KEY_SUNSET_FLAGGED' as any, + actor: creatorId, + creatorId: profile.id, + payload: { + ledger_sequence: Number(event.ledger), + flagged_at: flaggedAt.toISOString(), + }, + createdAt: flaggedAt, + }, + }), + ]); + + logger.info( + { + creatorId: profile.id, + sunsetFlaggedAt: flaggedAt.toISOString(), + ledger: event.ledger, + txHash: event.txHash, + }, + 'KEY_SUNSET_FLAGGED event processed; creator profile stamped' + ); + }); +} diff --git a/src/modules/indexer/trade-indexer.service.ts b/src/modules/indexer/trade-indexer.service.ts index 15bf0ce1..f5d5125f 100644 --- a/src/modules/indexer/trade-indexer.service.ts +++ b/src/modules/indexer/trade-indexer.service.ts @@ -9,6 +9,8 @@ export interface SorobanBuyEvent { ledger: number; tx_hash: string; timestamp: string; + /** Payment asset used for this purchase (e.g. 'XLM', 'USDC'). Optional; defaults to 'XLM' (#934). */ + payment_asset?: string; } const REQUIRED_FIELDS: (keyof SorobanBuyEvent)[] = [ @@ -84,6 +86,10 @@ export async function processTradeEvent( ledger: event.ledger, txHash: event.tx_hash, timestamp: new Date(event.timestamp), + paymentAsset: + typeof event.payment_asset === 'string' && event.payment_asset.trim() !== '' + ? event.payment_asset.trim().toUpperCase() + : 'XLM', }, }); diff --git a/src/modules/keys/key-analytics.service.ts b/src/modules/keys/key-analytics.service.ts index 856060fe..c5460318 100644 --- a/src/modules/keys/key-analytics.service.ts +++ b/src/modules/keys/key-analytics.service.ts @@ -54,6 +54,23 @@ export interface PlatformAnalytics extends TradeStats { to: string | null; } +export interface PaymentAssetStats { + payment_asset: string; + trade_count: number; + unique_traders: number; + total_volume: string; + trade_share_percent: number; +} + +export interface KeyPaymentAssetAnalytics { + keyId: string; + assets: PaymentAssetStats[]; +} + +export interface PlatformPaymentAssetDistribution { + assets: PaymentAssetStats[]; +} + function windowSuffix(window: AnalyticsWindow): string { return `${window.from?.toISOString() ?? '-'}:${window.to?.toISOString() ?? '-'}`; } @@ -78,7 +95,12 @@ export function getPlatformAnalyticsCacheKey( export async function invalidateKeyAnalyticsCache( keyId: string ): Promise { - await cacheInvalidate(`key:analytics:${keyId}:*`, 'platform:analytics:*'); + await cacheInvalidate( + `key:analytics:${keyId}:*`, + `key:payment-asset-analytics:${keyId}`, + 'platform:analytics:*', + 'platform:payment-asset-distribution' + ); } function buildTimestampFilter(window: AnalyticsWindow) { @@ -169,3 +191,80 @@ export async function getPlatformAnalytics( await cacheSetJson(cacheKey, analytics, KEY_ANALYTICS_CACHE_TTL_SECONDS); return analytics; } + +function aggregatePaymentAssets( + trades: Array<{ + buyer: string; + price: string; + quantity: string; + paymentAsset: string; + }> +): PaymentAssetStats[] { + const groups = new Map< + string, + { buyers: Set; tradeCount: number; volume: bigint } + >(); + + for (const trade of trades) { + const asset = trade.paymentAsset || 'XLM'; + const group = groups.get(asset) ?? { + buyers: new Set(), + tradeCount: 0, + volume: 0n, + }; + group.buyers.add(trade.buyer); + group.tradeCount += 1; + group.volume += BigInt(trade.price) * BigInt(trade.quantity); + groups.set(asset, group); + } + + return [...groups.entries()] + .map(([payment_asset, group]) => ({ + payment_asset, + trade_count: group.tradeCount, + unique_traders: group.buyers.size, + total_volume: group.volume.toString(), + trade_share_percent: trades.length + ? Number((BigInt(group.tradeCount) * 10000n) / BigInt(trades.length)) / + 100 + : 0, + })) + .sort((a, b) => a.payment_asset.localeCompare(b.payment_asset)); +} + +export async function getKeyPaymentAssetAnalytics( + keyId: string +): Promise { + const cacheKey = `key:payment-asset-analytics:${keyId}`; + const cached = await cacheGetJson(cacheKey); + if (cached) return cached; + + const creator = await prisma.creatorProfile.findUnique({ + where: { id: keyId }, + select: { id: true }, + }); + if (!creator) throw new KeyNotFoundError(keyId); + + const trades = await prisma.trade.findMany({ + where: { creatorId: creator.id }, + select: { buyer: true, price: true, quantity: true, paymentAsset: true }, + }); + const analytics = { keyId: creator.id, assets: aggregatePaymentAssets(trades) }; + + await cacheSetJson(cacheKey, analytics, KEY_ANALYTICS_CACHE_TTL_SECONDS); + return analytics; +} + +export async function getPlatformPaymentAssetDistribution(): Promise { + const cacheKey = 'platform:payment-asset-distribution'; + const cached = await cacheGetJson(cacheKey); + if (cached) return cached; + + const trades = await prisma.trade.findMany({ + select: { buyer: true, price: true, quantity: true, paymentAsset: true }, + }); + const distribution = { assets: aggregatePaymentAssets(trades) }; + + await cacheSetJson(cacheKey, distribution, KEY_ANALYTICS_CACHE_TTL_SECONDS); + return distribution; +} diff --git a/src/modules/keys/key-proposal-votes.service.test.ts b/src/modules/keys/key-proposal-votes.service.test.ts index 8c8198f3..fa505854 100644 --- a/src/modules/keys/key-proposal-votes.service.test.ts +++ b/src/modules/keys/key-proposal-votes.service.test.ts @@ -11,7 +11,8 @@ jest.mock('../../utils/prisma.utils', () => ({ create: jest.fn(), }, activity: { create: jest.fn() }, - keyOwnership: { findUnique: jest.fn() }, + keyOwnership: { findUnique: jest.fn(), aggregate: jest.fn() }, + voteDelegation: { findMany: jest.fn(), count: jest.fn() }, $transaction: jest.fn(), }, })); @@ -43,7 +44,8 @@ interface MockPrismaClient { }; governanceVote: { findUnique: jest.Mock; create: jest.Mock }; activity: { create: jest.Mock }; - keyOwnership: { findUnique: jest.Mock }; + keyOwnership: { findUnique: jest.Mock; aggregate: jest.Mock }; + voteDelegation: { findMany: jest.Mock; count: jest.Mock }; $transaction: jest.Mock; } @@ -98,6 +100,9 @@ function wireSupportingMocks( mockPrisma.keyOwnership.findUnique.mockResolvedValue({ balance: balance ?? 0, }); + mockPrisma.keyOwnership.aggregate.mockResolvedValue({ _sum: { balance: 0 } }); + mockPrisma.voteDelegation.findMany.mockResolvedValue([]); + mockPrisma.voteDelegation.count.mockResolvedValue(0); mockPrisma.activity.create.mockResolvedValue({}); mockPrisma.$transaction.mockImplementation(async (cb: any) => cb(mockPrisma)); } @@ -204,14 +209,39 @@ describe('castKeyProposalVote tallying', () => { expect(mockPrisma.activity.create).toHaveBeenCalledWith( expect.objectContaining({ data: expect.objectContaining({ - type: 'GOVERNANCE_PROPOSAL_CREATED', + type: 'GOVERNANCE_VOTE_CAST', actor: 'walletA', - payload: expect.objectContaining({ action: 'vote_cast' }), + payload: expect.objectContaining({ totalWeight: '5' }), }), }) ); }); + it('includes active delegators in the persisted vote weight', async () => { + const store = makeStore(); + wireProposalMocks(store); + wireSupportingMocks({ voted: false, balance: 5 }); + mockPrisma.voteDelegation.findMany.mockResolvedValue([ + { delegatorWallet: 'walletB' }, + { delegatorWallet: 'walletC' }, + ]); + mockPrisma.voteDelegation.count.mockResolvedValue(2); + mockPrisma.keyOwnership.aggregate.mockResolvedValue({ + _sum: { balance: 7 }, + }); + + const result = await castKeyProposalVote(KEY_ID, PROPOSAL_ID, 0, 'walletA'); + + expect(result).toMatchObject({ + weight: '12', + ownWeight: '5', + delegatedWeight: '7', + delegatorCount: 2, + }); + expect(store.total).toBe('12'); + expect(store.results).toEqual({ Yes: '12' }); + }); + it('accumulates tallies across multiple votes cast by different wallets on different options', async () => { const store = makeStore(); wireProposalMocks(store); diff --git a/src/modules/keys/key-proposal-votes.service.ts b/src/modules/keys/key-proposal-votes.service.ts index 105ddd95..1b7cbc26 100644 --- a/src/modules/keys/key-proposal-votes.service.ts +++ b/src/modules/keys/key-proposal-votes.service.ts @@ -2,6 +2,7 @@ import { prisma } from '../../utils/prisma.utils'; import { logger } from '../../utils/logger.utils'; import { Decimal } from '@prisma/client/runtime/library'; +import { getDelegatedVoteWeight } from '../governance/governance-delegation.service'; export class HolderNotEligibleError extends Error { constructor(wallet: string) { @@ -78,7 +79,14 @@ export interface CastVoteResult { proposalId: string; optionIndex: number; option: string; + /** Total vote weight including delegated weight from active delegators. */ weight: string; + /** The voter's own key balance. */ + ownWeight: string; + /** Sum of delegated weight from all active delegators. */ + delegatedWeight: string; + /** Number of active delegators whose weight was included. */ + delegatorCount: number; } /** @@ -128,9 +136,16 @@ export async function hasWalletVoted( /** * Submit a governance vote on behalf of a key holder and persist it. * - * The voter's key balance becomes the vote weight. The vote record is - * written to the `proposal_votes` table; a duplicate vote surfaces as a - * Prisma unique constraint violation mapped by the route to 409. + * Vote weight = voter's own key balance + sum of key balances of all wallets + * that have an active vote delegation pointing to this voter on the same key. + * Revoked delegations are excluded from the weight calculation. + * + * The vote record is written to the `proposal_votes` table with the combined + * weight. A duplicate vote surfaces as a Prisma unique constraint violation + * mapped by the route to 409. + * + * @returns CastVoteResult with the total weight broken down into ownWeight + * and delegatedWeight so callers can inspect the composition. */ export async function castKeyProposalVote( keyId: string, @@ -150,14 +165,15 @@ export async function castKeyProposalVote( throw new OptionIndexOutOfRangeError(optionIndex, options.length); } + // Own balance — determines eligibility. const ownership = await prisma.keyOwnership.findUnique({ where: { ownerAddress_creatorId: { ownerAddress: wallet, creatorId: keyId }, }, }); - const balance = ownership ? Number(ownership.balance) : 0; - if (balance <= 0) { + const ownBalance = ownership ? Number(ownership.balance) : 0; + if (ownBalance <= 0) { throw new HolderNotEligibleError(wallet); } @@ -166,7 +182,17 @@ export async function castKeyProposalVote( throw new DuplicateVoteError(); } - const weight = String(balance); + const [delegatedBalance, delegatorCount] = await Promise.all([ + getDelegatedVoteWeight(wallet, keyId), + prisma.voteDelegation.count({ + where: { delegateeWallet: wallet, keyId, isActive: true }, + }), + ]); + + const totalWeight = ownBalance + delegatedBalance; + const weightStr = String(totalWeight); + const ownWeightStr = String(ownBalance); + const delegatedWeightStr = String(delegatedBalance); const option = options[optionIndex]; // TODO: submit cast_vote contract call via Stellar SDK @@ -179,7 +205,10 @@ export async function castKeyProposalVote( voter: wallet, optionIndex, option, - weight, + ownWeight: ownWeightStr, + delegatedWeight: delegatedWeightStr, + totalWeight: weightStr, + delegatorCount, }, 'Submitting cast_vote contract call' ); @@ -210,7 +239,7 @@ export async function castKeyProposalVote( totalVotingWeight: proposal.totalVotingWeight, results: proposal.results as Record, }, - weight, + weightStr, option ); @@ -228,22 +257,24 @@ export async function castKeyProposalVote( proposalId, voter: wallet, optionIndex, - weight: new Decimal(weight), + weight: new Decimal(weightStr), }, }); await tx.activity.create({ data: { - type: 'GOVERNANCE_PROPOSAL_CREATED', + type: 'GOVERNANCE_VOTE_CAST' as any, actor: wallet, creatorId: keyId, payload: { keyId, proposalId, - action: 'vote_cast', optionIndex, option, - weight, + ownWeight: ownWeightStr, + delegatedWeight: delegatedWeightStr, + totalWeight: weightStr, + delegatorCount, }, }, }); @@ -253,6 +284,9 @@ export async function castKeyProposalVote( proposalId, optionIndex, option: options[optionIndex], - weight, + weight: weightStr, + ownWeight: ownWeightStr, + delegatedWeight: delegatedWeightStr, + delegatorCount, }; } diff --git a/src/modules/keys/key-sunset-watch.service.ts b/src/modules/keys/key-sunset-watch.service.ts new file mode 100644 index 00000000..349a2a06 --- /dev/null +++ b/src/modules/keys/key-sunset-watch.service.ts @@ -0,0 +1,202 @@ +// src/modules/keys/key-sunset-watch.service.ts +// Service backing GET /keys/sunset-watch (#931). +// +// Returns all creator keys that are: +// (a) already flagged on-chain (sunsetFlaggedAt IS NOT NULL), OR +// (b) approaching the inactivity sunset threshold — i.e. the most recent +// KEY_BOUGHT or KEY_SOLD activity for the key is older than +// KEY_SUNSET_INACTIVITY_THRESHOLD_DAYS days (or the key has never traded). +// +// Each result item carries: +// - daysSinceLastTrade — whole days since the last trade Activity record, +// or null when the key has never traded +// - sunsetStatus — 'sunset_pending' : on-chain flag is set +// 'threshold_exceeded' : last trade ≥ threshold days ago +// (or never traded) +// 'near_threshold' : last trade is within the +// threshold but past the +// near-window cutoff +// +// Results are sorted by daysSinceLastTrade descending (most inactive first); +// keys that have never traded sort last within the list. + +import { prisma } from '../../utils/prisma.utils'; +import { envConfig } from '../../config'; +import { buildOffsetPaginationMeta } from '../../utils/pagination.utils'; + +// ── Types ──────────────────────────────────────────────────── + +export type SunsetStatus = + | 'sunset_pending' + | 'threshold_exceeded' + | 'near_threshold'; + +export interface SunsetWatchItem { + keyId: string; + handle: string; + displayName: string; + circulatingSupply: string; + /** ISO timestamp of the last KEY_BOUGHT / KEY_SOLD activity, or null. */ + lastTradeAt: string | null; + /** Whole days elapsed since last trade, or null if never traded. */ + daysSinceLastTrade: number | null; + /** ISO timestamp when the on-chain KeySunsetFlagged event was processed, or null. */ + sunsetFlaggedAt: string | null; + sunsetStatus: SunsetStatus; +} + +export interface SunsetWatchResult { + items: SunsetWatchItem[]; + meta: ReturnType; +} + +export interface SunsetWatchQuery { + limit: number; + offset: number; +} + +// ── Internal helpers ───────────────────────────────────────── + +/** Whole elapsed days between `date` and now, floored to the nearest integer. */ +function daysSince(date: Date): number { + return Math.floor((Date.now() - date.getTime()) / 86_400_000); +} + +/** + * Derive the human-readable sunset status for a single key. + * + * Near-threshold window: the lesser of 7 days or half the configured threshold. + * A key is "near_threshold" when its last trade falls inside that window + * before the threshold, but has not been flagged on-chain. + */ +function deriveSunsetStatus( + sunsetFlaggedAt: Date | null, + daysSinceLastTrade: number | null, + thresholdDays: number, + nearWindowDays: number +): SunsetStatus { + if (sunsetFlaggedAt !== null) { + return 'sunset_pending'; + } + if ( + daysSinceLastTrade === null || + daysSinceLastTrade >= thresholdDays + ) { + return 'threshold_exceeded'; + } + // daysSinceLastTrade < thresholdDays, but within the near window + if (daysSinceLastTrade >= thresholdDays - nearWindowDays) { + return 'near_threshold'; + } + // Should not reach here (caller only passes candidates in the window), + // but guard defensively. + return 'near_threshold'; +} + +// ── Main query ─────────────────────────────────────────────── + +/** + * Fetch all keys approaching or past the inactivity sunset threshold plus + * all keys already flagged on-chain, sorted by inactivity descending. + * + * Strategy: + * 1. Pull all CreatorProfile rows that have sunsetFlaggedAt set. + * 2. Pull the max(createdAt) per creatorId from Activity for trade types, + * using Prisma groupBy, to find each key's last trade date. + * 3. Merge: include a key if it is flagged OR its last trade predates the + * near-threshold cutoff OR it has never traded. + * 4. Sort, then paginate in-memory (the eligible set is bounded and small). + */ +export async function getSunsetWatchList( + query: SunsetWatchQuery +): Promise { + const { limit, offset } = query; + const thresholdDays = envConfig.KEY_SUNSET_INACTIVITY_THRESHOLD_DAYS; + const nearWindowDays = Math.min(7, Math.floor(thresholdDays / 2)); + const nearCutoff = new Date( + Date.now() - (thresholdDays - nearWindowDays) * 86_400_000 + ); + + // ── 1. Fetch all profiles in a single query ────────────── + // We need all profiles to cross-reference against the activity aggregation. + // CreatorProfile counts are bounded (one per creator), so a full scan is fine. + const allProfiles = await prisma.creatorProfile.findMany({ + select: { + id: true, + handle: true, + displayName: true, + circulatingSupply: true, + sunsetFlaggedAt: true, + }, + }); + + // ── 2. Aggregate last trade date per creatorId ─────────── + // groupBy returns one row per creatorId with the latest createdAt. + const lastTradeRows = await prisma.activity.groupBy({ + by: ['creatorId'], + where: { + type: { in: ['KEY_BOUGHT', 'KEY_SOLD'] }, + creatorId: { not: null }, + }, + _max: { createdAt: true }, + }); + + const lastTradeMap = new Map(); + for (const row of lastTradeRows) { + if (row.creatorId && row._max.createdAt) { + lastTradeMap.set(row.creatorId, row._max.createdAt); + } + } + + // ── 3. Filter, enrich, and sort ────────────────────────── + type EnrichedItem = SunsetWatchItem & { _sortKey: number }; + const candidates: EnrichedItem[] = []; + + for (const profile of allProfiles) { + const lastTradeDate = lastTradeMap.get(profile.id) ?? null; + const days = lastTradeDate !== null ? daysSince(lastTradeDate) : null; + + // Include the key if it is flagged on-chain OR its last trade is old + // enough to be near/past the threshold OR it has never traded. + const isFlagged = profile.sunsetFlaggedAt !== null; + const isOldEnough = + lastTradeDate === null || lastTradeDate < nearCutoff; + + if (!isFlagged && !isOldEnough) { + continue; + } + + candidates.push({ + keyId: profile.id, + handle: profile.handle, + displayName: profile.displayName, + circulatingSupply: profile.circulatingSupply.toString(), + lastTradeAt: lastTradeDate?.toISOString() ?? null, + daysSinceLastTrade: days, + sunsetFlaggedAt: profile.sunsetFlaggedAt?.toISOString() ?? null, + sunsetStatus: deriveSunsetStatus( + profile.sunsetFlaggedAt, + days, + thresholdDays, + nearWindowDays + ), + // Never-traded keys sort after known-inactive keys (use large sentinel). + _sortKey: days ?? Number.MAX_SAFE_INTEGER, + }); + } + + // Sort descending: most inactive first; never-traded last. + candidates.sort((a, b) => b._sortKey - a._sortKey); + + // ── 4. Paginate ────────────────────────────────────────── + const total = candidates.length; + const page = candidates.slice(offset, offset + limit); + const items: SunsetWatchItem[] = page.map( + ({ _sortKey: _unused, ...item }) => item + ); + + return { + items, + meta: buildOffsetPaginationMeta({ limit, offset, total }), + }; +} diff --git a/src/modules/keys/keys.routes.ts b/src/modules/keys/keys.routes.ts index e0cc2b8e..2af70c59 100644 --- a/src/modules/keys/keys.routes.ts +++ b/src/modules/keys/keys.routes.ts @@ -105,6 +105,11 @@ import { calculateEffectiveWeight, } from '../staking/staking.service'; import { getKeyCurveMilestones } from './key-milestones.service'; +import { getSunsetWatchList } from './key-sunset-watch.service'; +import { + getKeyPaymentAssetAnalytics, + getPlatformPaymentAssetDistribution, +} from './key-analytics.service'; const priceHistoryQuerySchema = z.object({ from: z.string().datetime(), @@ -327,6 +332,128 @@ router.get('/search', async (req, res, next) => { } }); +// ── Pagination constants ──────────────────────────────────── +const DEFAULT_SUNSET_WATCH_LIMIT = 20; +const MAX_SUNSET_WATCH_LIMIT = 100; + +const sunsetWatchQuerySchema = z.object({ + limit: z.coerce + .number() + .int() + .min(1) + .max(MAX_SUNSET_WATCH_LIMIT) + .default(DEFAULT_SUNSET_WATCH_LIMIT), + offset: z.coerce.number().int().min(0).default(0), +}); + +/** + * GET /api/v1/keys/sunset-watch + * + * Admin-only endpoint that returns all creator keys approaching or past the + * inactivity sunset threshold, plus any keys already flagged on-chain via a + * KeySunsetFlagged event. + * + * Each item includes: + * - keyId, handle, displayName, circulatingSupply + * - lastTradeAt — ISO timestamp of the last KEY_BOUGHT/KEY_SOLD, or null + * - daysSinceLastTrade — whole days elapsed since last trade, or null + * - sunsetFlaggedAt — ISO timestamp the on-chain flag was processed, or null + * - sunsetStatus — 'sunset_pending' | 'threshold_exceeded' | 'near_threshold' + * + * Results are sorted by inactivity duration descending (most inactive first). + * Keys that have never traded appear after all keys with a known last trade. + * + * Query parameters: + * - limit (default 20, max 100) + * - offset (default 0) + * + * Responses: + * 200 — paginated list of sunset-watch items + * 400 — invalid query parameters + * 401 — missing or invalid admin token + * 403 — token present but role !== 'admin' + * + * Must be registered before /:keyId to avoid route shadowing. + */ +router.get( + '/sunset-watch', + adminGuard, + async (req: AdminRequest, res, next) => { + const parsed = sunsetWatchQuerySchema.safeParse(req.query); + if (!parsed.success) { + sendValidationError( + res, + 'Invalid query parameters', + zodIssuesToDetails(parsed.error.issues) + ); + return; + } + + try { + const result = await getSunsetWatchList({ + limit: parsed.data.limit, + offset: parsed.data.offset, + }); + sendSuccess(res, result); + } catch (error) { + logger.error({ error }, 'GET /keys/sunset-watch failed'); + next(error); + } + } +); + +/** + * GET /api/v1/keys/analytics/payment-assets + * + * Admin-only. Returns the platform-wide distribution of payment assets used + * across all key purchases — trade count, unique buyers, total price in + * stroops, and percentage share per asset. + * + * Must be registered before /:keyId to avoid route shadowing. + * + * Responses: + * 200 — PlatformPaymentAssetDistribution + * 401/403 — missing or invalid admin token + */ +router.get( + '/analytics/payment-assets', + adminGuard, + async (_req: AdminRequest, res, next) => { + try { + sendSuccess(res, await getPlatformPaymentAssetDistribution()); + } catch (error) { + logger.error({ error }, 'GET /keys/analytics/payment-assets failed'); + next(error); + } + } +); + +/** + * GET /api/v1/keys/:keyId/analytics + * + * Returns the payment-asset breakdown for a single key: trade count, unique + * buyers, and total price in stroops grouped by paymentAsset. + * Resolves keyId by DB id or handle. + * + * Publicly accessible — the payment-asset mix for a key is not sensitive. + * + * Responses: + * 200 — KeyPaymentAssetAnalytics + * 404 — key not found + */ +router.get('/:keyId/analytics', async (req, res, next) => { + try { + sendSuccess(res, await getKeyPaymentAssetAnalytics(String(req.params.keyId))); + } catch (error) { + if (error instanceof KeyNotFoundError) { + sendNotFound(res, 'Key'); + return; + } + logger.error({ error, keyId: req.params.keyId }, 'GET /keys/:keyId/analytics failed'); + next(error); + } +}); + /** * GET /api/v1/keys/:keyId/oracle-price * diff --git a/src/modules/staking/staking-indexer.service.ts b/src/modules/staking/staking-indexer.service.ts new file mode 100644 index 00000000..d1a4b928 --- /dev/null +++ b/src/modules/staking/staking-indexer.service.ts @@ -0,0 +1,356 @@ +// src/modules/staking/staking-indexer.service.ts +// Indexer event handlers for staking NFT contract events (#932). +// +// Two event types are handled: +// +// STAKE_NFT_MINTED +// Emitted when the staking contract mints a new stake receipt NFT. +// Creates a StakingNft row, a StakingNftTransfer row (fromAddress=null), +// and an Activity record for the audit trail. Idempotent via mintTxHash. +// +// STAKE_NFT_TRANSFERRED +// Emitted when a stake receipt NFT is transferred from one wallet to +// another (including secondary-market sales). Updates ownerAddress on the +// StakingNft row and appends a StakingNftTransfer row. Idempotent via +// the (txHash, eventIndex) unique constraint on StakingNftTransfer. +// +// Both handlers invalidate the affected wallets' NFT list caches so that the +// read endpoints reflect the change within the next request cycle. + +import { prisma } from '../../utils/prisma.utils'; +import { logger } from '../../utils/logger.utils'; +import { + processIndexerChainEvents, + IndexerChainEvent, +} from '../../utils/indexer-event-processor.utils'; +import { cacheInvalidate } from '../../utils/redis.utils'; +import { + walletNftsCachePattern, + nftMetaCachePattern, + nftTransfersCachePattern, +} from './staking.service'; + +// ── Typed event interfaces ──────────────────────────────────── + +/** + * Contract event emitted when a new staking NFT is minted. + * + * Required fields (all validated before any DB write): + * tokenId — on-chain token identifier (string) + * ownerAddress — wallet that receives the NFT on mint + * keyId — creator key the staked tokens belong to + * stakedAmount — number of locked keys (numeric string) + * mintTxHash — transaction hash of the mint (dedup key) + * ledger — ledger sequence of the mint + * mintedAt — ISO-8601 timestamp of the mint + * + * Optional fields: + * lockExpiryLedger — ledger at which the lock expires (omit if no lock) + * lockExpiresAt — ISO-8601 wall-clock equivalent of lockExpiryLedger + */ +export interface StakeNftMintedEvent extends IndexerChainEvent { + eventType: 'STAKE_NFT_MINTED'; + tokenId: string; + ownerAddress: string; + keyId: string; + stakedAmount: string; + mintTxHash: string; + mintedAt: string; + lockExpiryLedger?: number; + lockExpiresAt?: string; +} + +/** + * Contract event emitted when a staking NFT changes owner. + * + * Required fields: + * tokenId — identifies which NFT was transferred + * fromAddress — previous owner + * toAddress — new owner + * ledger — ledger sequence of the transfer + * txHash — transaction hash + * eventIndex — position within the transaction (dedup with txHash) + * occurredAt — ISO-8601 timestamp of the transfer + */ +export interface StakeNftTransferredEvent extends IndexerChainEvent { + eventType: 'STAKE_NFT_TRANSFERRED'; + tokenId: string; + fromAddress: string; + toAddress: string; + occurredAt: string; +} + +// ── STAKE_NFT_MINTED ────────────────────────────────────────── + +const MINT_REQUIRED_FIELDS: (keyof StakeNftMintedEvent)[] = [ + 'tokenId', + 'ownerAddress', + 'keyId', + 'stakedAmount', + 'mintTxHash', + 'mintedAt', + 'ledger', +]; + +/** + * Process a batch of STAKE_NFT_MINTED events. + * + * Each event is written in a single Prisma transaction that creates: + * 1. The StakingNft row + * 2. An initial StakingNftTransfer row (fromAddress = null, representing + * the mint itself) + * 3. An Activity record for the audit trail + * + * Already-minted NFTs (matched by mintTxHash) are skipped so replays are + * safe. + */ +export async function processStakeNftMintedEvents( + events: IndexerChainEvent[] +): Promise { + await processIndexerChainEvents(events, async event => { + if (event.eventType !== 'STAKE_NFT_MINTED') return; + + const e = event as StakeNftMintedEvent; + + for (const field of MINT_REQUIRED_FIELDS) { + const value = e[field]; + if (value === undefined || value === null || value === '') { + logger.warn( + { + eventId: `${e.txHash}:${e.eventIndex}`, + missingField: field, + }, + 'Skipping STAKE_NFT_MINTED event due to missing required field' + ); + return; + } + } + + // Idempotency: skip if we already processed this mint transaction. + const existing = await prisma.stakingNft.findUnique({ + where: { mintTxHash: e.mintTxHash }, + select: { id: true }, + }); + if (existing) { + logger.info( + { + eventId: `${e.txHash}:${e.eventIndex}`, + tokenId: e.tokenId, + mintTxHash: e.mintTxHash, + }, + 'STAKE_NFT_MINTED already recorded; skipping duplicate event' + ); + return; + } + + const mintedAt = new Date(e.mintedAt); + const lockExpiresAt = + e.lockExpiresAt ? new Date(e.lockExpiresAt) : null; + + const nft = await prisma.$transaction(async tx => { + // 1. Create the NFT record. + const created = await tx.stakingNft.create({ + data: { + tokenId: e.tokenId, + ownerAddress: e.ownerAddress, + keyId: e.keyId, + stakedAmount: e.stakedAmount, + lockExpiryLedger: e.lockExpiryLedger ?? null, + lockExpiresAt, + mintLedger: Number(e.ledger), + mintTxHash: e.mintTxHash, + }, + }); + + // 2. Record the initial "mint" transfer (fromAddress = null). + await tx.stakingNftTransfer.create({ + data: { + nftId: created.id, + fromAddress: null, + toAddress: e.ownerAddress, + ledger: Number(e.ledger), + txHash: e.txHash, + eventIndex: Number(e.eventIndex), + occurredAt: mintedAt, + }, + }); + + // 3. Audit trail. + await tx.activity.create({ + data: { + type: 'STAKE_NFT_MINTED' as any, + actor: e.ownerAddress, + creatorId: e.keyId, + payload: { + tokenId: e.tokenId, + stakedAmount: e.stakedAmount, + lockExpiryLedger: e.lockExpiryLedger ?? null, + ledger_sequence: Number(e.ledger), + }, + createdAt: mintedAt, + }, + }); + + return created; + }); + + // Invalidate wallet cache so the new NFT appears immediately. + await cacheInvalidate(walletNftsCachePattern(e.ownerAddress)); + + logger.info( + { + nftId: nft.id, + tokenId: e.tokenId, + ownerAddress: e.ownerAddress, + keyId: e.keyId, + ledger: e.ledger, + txHash: e.txHash, + }, + 'STAKE_NFT_MINTED event processed' + ); + }); +} + +// ── STAKE_NFT_TRANSFERRED ───────────────────────────────────── + +const TRANSFER_REQUIRED_FIELDS: (keyof StakeNftTransferredEvent)[] = [ + 'tokenId', + 'fromAddress', + 'toAddress', + 'occurredAt', + 'ledger', + 'txHash', + 'eventIndex', +]; + +/** + * Process a batch of STAKE_NFT_TRANSFERRED events. + * + * For each event: + * 1. Looks up the StakingNft by tokenId — skips with a warning if not found + * (out-of-order delivery; a replay from a full resync will fix it). + * 2. Updates ownerAddress on the StakingNft row. + * 3. Appends a StakingNftTransfer history row. + * 4. Writes an Activity record. + * + * Idempotent via the unique(txHash, eventIndex) constraint on + * StakingNftTransfer — a duplicate write attempt is silently skipped. + */ +export async function processStakeNftTransferredEvents( + events: IndexerChainEvent[] +): Promise { + await processIndexerChainEvents(events, async event => { + if (event.eventType !== 'STAKE_NFT_TRANSFERRED') return; + + const e = event as StakeNftTransferredEvent; + + for (const field of TRANSFER_REQUIRED_FIELDS) { + const value = e[field]; + if (value === undefined || value === null || value === '') { + logger.warn( + { + eventId: `${e.txHash}:${e.eventIndex}`, + missingField: field, + }, + 'Skipping STAKE_NFT_TRANSFERRED event due to missing required field' + ); + return; + } + } + + // Idempotency: skip if this exact transfer is already recorded. + const existingTransfer = await prisma.stakingNftTransfer.findUnique({ + where: { + txHash_eventIndex: { + txHash: e.txHash, + eventIndex: Number(e.eventIndex), + }, + }, + select: { id: true }, + }); + if (existingTransfer) { + logger.info( + { eventId: `${e.txHash}:${e.eventIndex}`, tokenId: e.tokenId }, + 'STAKE_NFT_TRANSFERRED already recorded; skipping duplicate event' + ); + return; + } + + // Resolve the NFT by on-chain tokenId. + const nft = await prisma.stakingNft.findUnique({ + where: { tokenId: e.tokenId }, + select: { id: true, ownerAddress: true }, + }); + + if (!nft) { + logger.warn( + { + eventId: `${e.txHash}:${e.eventIndex}`, + tokenId: e.tokenId, + }, + 'STAKE_NFT_TRANSFERRED references unknown tokenId; skipping (will retry on resync)' + ); + return; + } + + const occurredAt = new Date(e.occurredAt); + const previousOwner = nft.ownerAddress; + + await prisma.$transaction([ + // Update the current owner. + prisma.stakingNft.update({ + where: { id: nft.id }, + data: { ownerAddress: e.toAddress }, + }), + // Append a transfer history row. + prisma.stakingNftTransfer.create({ + data: { + nftId: nft.id, + fromAddress: e.fromAddress, + toAddress: e.toAddress, + ledger: Number(e.ledger), + txHash: e.txHash, + eventIndex: Number(e.eventIndex), + occurredAt, + }, + }), + // Audit trail. + prisma.activity.create({ + data: { + type: 'STAKE_NFT_TRANSFERRED' as any, + actor: e.fromAddress, + target: e.toAddress, + creatorId: null, + payload: { + tokenId: e.tokenId, + nftId: nft.id, + ledger_sequence: Number(e.ledger), + }, + createdAt: occurredAt, + }, + }), + ]); + + // Invalidate both wallets' NFT list caches and the per-NFT caches. + await cacheInvalidate( + walletNftsCachePattern(previousOwner), + walletNftsCachePattern(e.toAddress), + nftMetaCachePattern(nft.id), + nftMetaCachePattern(e.tokenId), + nftTransfersCachePattern(nft.id), + nftTransfersCachePattern(e.tokenId) + ); + + logger.info( + { + nftId: nft.id, + tokenId: e.tokenId, + fromAddress: e.fromAddress, + toAddress: e.toAddress, + ledger: e.ledger, + txHash: e.txHash, + }, + 'STAKE_NFT_TRANSFERRED event processed' + ); + }); +} diff --git a/src/modules/staking/staking.service.ts b/src/modules/staking/staking.service.ts index 4fe5b349..a4ca458c 100644 --- a/src/modules/staking/staking.service.ts +++ b/src/modules/staking/staking.service.ts @@ -483,3 +483,19 @@ export async function getPositionEffectiveWeight( tier: position.tierData, }; } + +// ── NFT cache-key helpers (#932) ───────────────────────────── +// These are consumed by staking-indexer.service.ts to invalidate the Redis +// caches used by the staking NFT read endpoints after mint/transfer events. + +export function walletNftsCachePattern(wallet: string): string { + return `staking:nfts:wallet:${wallet}:*`; +} + +export function nftMetaCachePattern(identifier: string): string { + return `staking:nfts:meta:${identifier}`; +} + +export function nftTransfersCachePattern(identifier: string): string { + return `staking:nfts:transfers:${identifier}:*`; +}