From cbef51543703dfc82f57ee606e4f93bf908f95d5 Mon Sep 17 00:00:00 2001 From: Harkinkunmi <158626915+TheHalalHunter@users.noreply.github.com> Date: Sat, 26 Sep 2026 19:32:56 +0100 Subject: [PATCH 1/5] feat: implement issues #931-#934 #931 Add sunset tracking endpoint for keys approaching inactivity threshold - sunsetFlaggedAt on CreatorProfile + migration - KEY_SUNSET_FLAGGED ActivityType value - key-sunset-watch.service.ts (inactivity + near-threshold logic) - processSunsetFlaggedEvents() indexer handler - GET /keys/sunset-watch (admin-only, sorted by inactivity desc) - KEY_SUNSET_INACTIVITY_THRESHOLD_DAYS env var (default 30) #932 Build staking NFT position endpoints for tradeable stake receipt tracking - StakingNft + StakingNftTransfer schema + migration - STAKE_NFT_MINTED/TRANSFERRED ActivityType values - staking.service.ts read-model with Redis caching - staking-indexer.service.ts (processStakeNftMintedEvents, processStakeNftTransferredEvents) - GET /staking/nfts (auth), GET /staking/nfts/:id (public), GET /staking/nfts/:id/transfers (public) #933 Implement delegated governance voting endpoints with delegation tracking - VoteDelegation + VoteDelegationHistory schema + migration - GOVERNANCE_VOTE_CAST/DELEGATION_SET/DELEGATION_REVOKED ActivityType values - governance-delegation.service.ts read-model - governance-delegation-indexer.service.ts (processDelegationSetEvents, processDelegationRevokedEvents) - castKeyProposalVote updated: weight = own balance + sum of active delegators balances - GET /governance/delegation/:wallet, GET /governance/delegation/:wallet/history - GET /governance/delegators/:wallet #934 Add multi-currency payment support tracking for key purchase analytics - paymentAsset String @default("XLM") on Trade model + migration + indexes - payment_asset threaded through indexer-pipeline.service.ts Activity payload - payment_asset threaded through trade-indexer.service.ts db.trade.create - key-analytics.service.ts (getKeyPaymentAssetAnalytics, getPlatformPaymentAssetDistribution) - GET /keys/:keyId/analytics (public), GET /keys/analytics/payment-assets (admin) --- .env.example | 4 + prisma/schema/activity.prisma | 12 + prisma/schema/creator.prisma | 3 + prisma/schema/governance-delegation.prisma | 80 ++++ .../migration.sql | 12 + .../migration.sql | 69 +++ .../migration.sql | 71 +++ .../migration.sql | 11 + prisma/schema/staking.prisma | 93 ++++ prisma/schema/trade.prisma | 23 +- src/config.schema.ts | 10 + .../governance-delegation-indexer.service.ts | 438 ++++++++++++++++++ .../governance-delegation.service.ts | 347 ++++++++++++++ src/modules/governance/governance.routes.ts | 272 +++++++++++ src/modules/index.ts | 5 + .../indexer/indexer-pipeline.service.ts | 115 +++++ src/modules/indexer/trade-indexer.service.ts | 6 + src/modules/keys/key-analytics.service.ts | 212 +++++++++ .../keys/key-proposal-votes.service.ts | 62 ++- src/modules/keys/key-sunset-watch.service.ts | 202 ++++++++ src/modules/keys/keys.routes.ts | 127 +++++ .../staking/staking-indexer.service.ts | 356 ++++++++++++++ src/modules/staking/staking.routes.ts | 216 +++++++++ src/modules/staking/staking.service.ts | 278 +++++++++++ 24 files changed, 3003 insertions(+), 21 deletions(-) create mode 100644 prisma/schema/governance-delegation.prisma create mode 100644 prisma/schema/migrations/20260925000100_add_key_sunset_flagged_at/migration.sql create mode 100644 prisma/schema/migrations/20260925000200_add_staking_nfts/migration.sql create mode 100644 prisma/schema/migrations/20260925000300_add_governance_delegation/migration.sql create mode 100644 prisma/schema/migrations/20260925000400_add_trade_payment_asset/migration.sql create mode 100644 prisma/schema/staking.prisma create mode 100644 src/modules/governance/governance-delegation-indexer.service.ts create mode 100644 src/modules/governance/governance-delegation.service.ts create mode 100644 src/modules/governance/governance.routes.ts create mode 100644 src/modules/keys/key-analytics.service.ts create mode 100644 src/modules/keys/key-sunset-watch.service.ts create mode 100644 src/modules/staking/staking-indexer.service.ts create mode 100644 src/modules/staking/staking.routes.ts create mode 100644 src/modules/staking/staking.service.ts diff --git a/.env.example b/.env.example index e66a53e8..df37b49b 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/activity.prisma b/prisma/schema/activity.prisma index aaa9a6b4..32172de4 100644 --- a/prisma/schema/activity.prisma +++ b/prisma/schema/activity.prisma @@ -17,6 +17,18 @@ enum ActivityType { GOVERNANCE_PROPOSAL_CREATED TIMELOCK_LOCKUP_PROPOSED CIRCUIT_BREAKER_THRESHOLD_UPDATED + /// On-chain event emitted when a key is flagged for inactivity sunset (#931). + KEY_SUNSET_FLAGGED + /// On-chain event emitted when a new staking NFT is minted (#932). + STAKE_NFT_MINTED + /// On-chain event emitted when a staking NFT is transferred to a new holder (#932). + STAKE_NFT_TRANSFERRED + /// A governance vote was cast, optionally with delegated weight (#933). + GOVERNANCE_VOTE_CAST + /// A wallet set or updated its vote delegate (#933). + GOVERNANCE_DELEGATION_SET + /// A wallet revoked its active vote delegation (#933). + GOVERNANCE_DELEGATION_REVOKED } model Activity { diff --git a/prisma/schema/creator.prisma b/prisma/schema/creator.prisma index 8bb45df5..308585a6 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 new file mode 100644 index 00000000..aee2f814 --- /dev/null +++ b/prisma/schema/staking.prisma @@ -0,0 +1,93 @@ +// prisma/schema/staking.prisma +// Staking NFT position read model (#932). +// +// Each StakingNft represents a tradeable stake receipt minted on-chain when a +// wallet locks creator keys into the staking contract. The NFT carries the +// staked amount and lock expiry so it can be traded independently of the +// underlying key position. +// +// Ownership is kept current by processing StakeNFTMinted and +// StakeNFTTransferred contract events in the indexer pipeline (see +// src/modules/staking/staking-indexer.service.ts). +// +// StakingNftTransfer is the append-only ownership history log for a single +// NFT; every transfer event (including the initial mint-to-owner) is written +// here so GET /staking/nfts/:id/transfers can return the full provenance. + +model StakingNft { + id String @id @default(cuid()) + + /// On-chain token ID emitted by the contract (unique per network). + tokenId String @unique + + /// Wallet currently holding this NFT. + ownerAddress String + + /// Creator key the staked tokens belong to. + keyId String + + /// Number of creator keys locked in this NFT position (stored as string + /// to preserve full precision; convert to Decimal/BigInt in application code). + stakedAmount Decimal @db.Decimal(30, 7) + + /// Ledger at which the lock expires (null = no lock / already expired). + lockExpiryLedger Int? + + /// Wall-clock time at which the lock expires, derived from lockExpiryLedger + /// by the indexer (null until computed or if no lock). + lockExpiresAt DateTime? + + /// Whether this NFT has been burned / redeemed. Burned NFTs are kept for + /// history but excluded from active-position queries. + burned Boolean @default(false) + burnedAt DateTime? + + /// Ledger sequence of the minting transaction. + mintLedger Int + + /// Transaction hash of the minting transaction (used for idempotency). + 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. +/// The first row for each NFT is the mint event (fromAddress IS NULL). +model StakingNftTransfer { + id String @id @default(cuid()) + + nftId String + + /// Sender wallet — null on the initial mint. + fromAddress String? + + /// Recipient wallet. + toAddress String + + /// Ledger sequence of the transfer transaction. + ledger Int + + /// On-chain transaction hash (used for idempotency). + txHash String + + /// Index within the transaction (used together with txHash for idempotency). + 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 638429ff..a0079aea 100644 --- a/src/config.schema.ts +++ b/src/config.schema.ts @@ -228,6 +228,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/governance/governance.routes.ts b/src/modules/governance/governance.routes.ts new file mode 100644 index 00000000..ee4175db --- /dev/null +++ b/src/modules/governance/governance.routes.ts @@ -0,0 +1,272 @@ +// src/modules/governance/governance.routes.ts +// HTTP route handlers for delegated governance voting (#933). +// +// Endpoints +// ───────────────────────────────────────────────────────────────────────────── +// +// GET /governance/delegation/:wallet +// Public. Returns the current active delegate for a wallet. +// Optional query param ?keyId= narrows to a specific creator key. +// Returns { delegated: false } when no active delegation exists rather +// than 404, so the frontend can display "not delegated" cleanly. +// +// GET /governance/delegators/:wallet +// Public. Returns all wallets that have delegated to this address (paginated). +// Optional query param ?keyId= narrows to a specific creator key. +// +// GET /governance/delegation/:wallet/history +// Public. Returns the full event log of delegations set/revoked by or to +// this wallet, newest first (paginated). +// Optional query param ?keyId= narrows to a specific creator key. +// +// Route ordering: /:wallet/history must be registered BEFORE /:wallet so +// Express does not match "history" as the :wallet param value. + +import { Router } from 'express'; +import { z } from 'zod'; +import { + sendNotFound, + sendSuccess, + sendValidationError, + zodIssuesToDetails, +} from '../../utils/api-response.utils'; +import { logger } from '../../utils/logger.utils'; +import { StellarAddressSchema } from '../wallet/wallet.schemas'; +import { + getCurrentDelegate, + getActiveDelegators, + getDelegationHistory, +} from './governance-delegation.service'; + +const governanceRouter = Router(); + +// ── Shared constants ────────────────────────────────────────── + +const DEFAULT_LIMIT = 20; +const MAX_LIMIT = 100; + +// ── Shared Zod schemas ──────────────────────────────────────── + +/** Optional keyId filter — any non-empty string is valid. */ +const keyIdQuerySchema = z.object({ + keyId: z.string().min(1).optional(), +}); + +const paginationQuerySchema = z.object({ + limit: z.coerce.number().int().min(1).max(MAX_LIMIT).default(DEFAULT_LIMIT), + offset: z.coerce.number().int().min(0).default(0), + keyId: z.string().min(1).optional(), +}); + +const walletParamSchema = z.object({ + wallet: StellarAddressSchema, +}); + +// ── GET /governance/delegation/:wallet/history ──────────────── +// Registered BEFORE /:wallet to avoid "history" matching as the :wallet param. + +/** + * GET /api/v1/governance/delegation/:wallet/history + * + * Returns the full event log (DelegationSet and DelegationRevoked) for a + * wallet — both events where the wallet was the delegator and events where + * it was the delegatee — newest first. + * + * This endpoint is publicly accessible: delegation history is a governance + * transparency feature. + * + * Path parameters: + * wallet — Stellar address of the wallet + * + * Query parameters: + * keyId — (optional) narrow to a specific creator key + * limit — max items per page (1–100, default 20) + * offset — items to skip (default 0) + * + * Responses: + * 200 — { items: DelegationHistoryItem[], meta: OffsetPaginationMeta } + * 400 — invalid wallet address or query parameters + */ +governanceRouter.get('/:wallet/history', async (req, res, next) => { + const walletParsed = walletParamSchema.safeParse(req.params); + if (!walletParsed.success) { + sendValidationError( + res, + 'Invalid wallet address', + zodIssuesToDetails(walletParsed.error.issues) + ); + return; + } + + const queryParsed = paginationQuerySchema.safeParse(req.query); + if (!queryParsed.success) { + sendValidationError( + res, + 'Invalid query parameters', + zodIssuesToDetails(queryParsed.error.issues) + ); + return; + } + + try { + const result = await getDelegationHistory({ + wallet: walletParsed.data.wallet, + keyId: queryParsed.data.keyId, + limit: queryParsed.data.limit, + offset: queryParsed.data.offset, + }); + sendSuccess(res, result); + } catch (error) { + logger.error( + { error, wallet: req.params.wallet }, + 'GET /governance/delegation/:wallet/history failed' + ); + next(error); + } +}); + +governanceRouter.all('/:wallet/history', (_req, res) => { + res.set('Allow', 'GET').sendStatus(405); +}); + +// ── GET /governance/delegation/:wallet ──────────────────────── + +/** + * GET /api/v1/governance/delegation/:wallet + * + * Returns the current active delegate for a wallet. When the wallet has no + * active delegation, returns `{ delegated: false }` rather than 404 so the + * frontend can render the "not delegated" state without error handling. + * + * Optionally scope to a specific creator key with `?keyId=`. Without it, + * the most-recently-set active delegation is returned. + * + * Publicly accessible. + * + * Path parameters: + * wallet — Stellar address of the delegating wallet + * + * Query parameters: + * keyId — (optional) narrow to a specific creator key + * + * Responses: + * 200 — { delegated: false } | { delegated: true, delegation: DelegationItem } + * 400 — invalid wallet address or query parameters + */ +governanceRouter.get('/:wallet', async (req, res, next) => { + const walletParsed = walletParamSchema.safeParse(req.params); + if (!walletParsed.success) { + sendValidationError( + res, + 'Invalid wallet address', + zodIssuesToDetails(walletParsed.error.issues) + ); + return; + } + + const queryParsed = keyIdQuerySchema.safeParse(req.query); + if (!queryParsed.success) { + sendValidationError( + res, + 'Invalid query parameters', + zodIssuesToDetails(queryParsed.error.issues) + ); + return; + } + + try { + const result = await getCurrentDelegate({ + wallet: walletParsed.data.wallet, + keyId: queryParsed.data.keyId, + }); + sendSuccess(res, result); + } catch (error) { + logger.error( + { error, wallet: req.params.wallet }, + 'GET /governance/delegation/:wallet failed' + ); + next(error); + } +}); + +governanceRouter.all('/:wallet', (_req, res) => { + res.set('Allow', 'GET').sendStatus(405); +}); + +export { governanceRouter as delegationRouter }; + +// ── /governance/delegators sub-router ──────────────────────── +// Mounted separately in index.ts under /governance/delegators so the +// path shape matches the spec: GET /governance/delegators/:wallet + +const delegatorsRouter = Router(); + +/** + * GET /api/v1/governance/delegators/:wallet + * + * Returns all wallets that currently have an active delegation pointing to + * this address, sorted by delegation date descending. + * + * Useful for showing a delegate how much aggregate vote weight they carry + * and from whom. + * + * Optionally filter by creator key with `?keyId=`. + * + * Publicly accessible. + * + * Path parameters: + * wallet — Stellar address of the delegate (the recipient of delegations) + * + * Query parameters: + * keyId — (optional) narrow to a specific creator key + * limit — max items per page (1–100, default 20) + * offset — items to skip (default 0) + * + * Responses: + * 200 — { items: DelegatorItem[], meta: OffsetPaginationMeta } + * 400 — invalid wallet address or query parameters + */ +delegatorsRouter.get('/:wallet', async (req, res, next) => { + const walletParsed = walletParamSchema.safeParse(req.params); + if (!walletParsed.success) { + sendValidationError( + res, + 'Invalid wallet address', + zodIssuesToDetails(walletParsed.error.issues) + ); + return; + } + + const queryParsed = paginationQuerySchema.safeParse(req.query); + if (!queryParsed.success) { + sendValidationError( + res, + 'Invalid query parameters', + zodIssuesToDetails(queryParsed.error.issues) + ); + return; + } + + try { + const result = await getActiveDelegators({ + wallet: walletParsed.data.wallet, + keyId: queryParsed.data.keyId, + limit: queryParsed.data.limit, + offset: queryParsed.data.offset, + }); + sendSuccess(res, result); + } catch (error) { + logger.error( + { error, wallet: req.params.wallet }, + 'GET /governance/delegators/:wallet failed' + ); + next(error); + } +}); + +delegatorsRouter.all('/:wallet', (_req, res) => { + res.set('Allow', 'GET').sendStatus(405); +}); + +export { delegatorsRouter }; +export default governanceRouter; diff --git a/src/modules/index.ts b/src/modules/index.ts index 17136645..37e4dd7e 100644 --- a/src/modules/index.ts +++ b/src/modules/index.ts @@ -27,6 +27,8 @@ import protocolRouter from './protocol/protocol.routes'; import revenueRouter from './revenue/revenue.routes'; import stakerRouter from './revenue/staker-revenue.routes'; import portfolioRouter from './portfolio/portfolio.routes'; +import stakingRouter from './staking/staking.routes'; +import { delegationRouter, delegatorsRouter } from './governance/governance.routes'; import { BASE as CREATORS_BASE } from '../constants/creator.constants'; const router = Router(); @@ -71,5 +73,8 @@ router.use('/protocol', routeBodySizeLimit('default'), protocolRouter); router.use('/revenue', routeBodySizeLimit('default'), revenueRouter); router.use('/staker', routeBodySizeLimit('default'), stakerRouter); router.use('/portfolio', routeBodySizeLimit('default'), portfolioRouter); +router.use('/staking', routeBodySizeLimit('default'), stakingRouter); +router.use('/governance/delegation', routeBodySizeLimit('default'), delegationRouter); +router.use('/governance/delegators', routeBodySizeLimit('default'), delegatorsRouter); export default router; diff --git a/src/modules/indexer/indexer-pipeline.service.ts b/src/modules/indexer/indexer-pipeline.service.ts index 014768a5..900d243b 100644 --- a/src/modules/indexer/indexer-pipeline.service.ts +++ b/src/modules/indexer/indexer-pipeline.service.ts @@ -44,6 +44,11 @@ export async function processTradeEvents(events: IndexerChainEvent[]): Promise) .join('|'); return createHash('sha256').update(identifiers, 'utf8').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 0768c2f4..332c0fde 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 new file mode 100644 index 00000000..5ebd2611 --- /dev/null +++ b/src/modules/keys/key-analytics.service.ts @@ -0,0 +1,212 @@ +// src/modules/keys/key-analytics.service.ts +// Payment-asset analytics for key purchase tracking (#934). +// +// Two query surfaces: +// +// getKeyPaymentAssetAnalytics(keyId) +// Per-key breakdown: trade count, unique buyers, and total price (stroops) +// grouped by paymentAsset. Used by GET /keys/:keyId/analytics. +// +// getPlatformPaymentAssetDistribution() +// Platform-wide: same aggregation across ALL keys plus a percentage share +// per asset. Used by GET /keys/analytics/payment-assets (admin). +// +// Both surfaces read from the Trade table directly — it is the canonical +// source of truth for payment asset data. Activity.payload.payment_asset is +// written in parallel by the indexer (for event-stream consumers) but is not +// queried here to avoid parsing Json. +// +// Caching: both results are cached in Redis for 5 minutes. Cache is +// invalidated by keyId after each successful trade write (the trade-indexer +// already calls invalidateCreatorDashboardCache; callers should also call +// invalidateKeyAnalyticsCache). + +import { prisma } from '../../utils/prisma.utils'; +import { cacheGetJson, cacheSetJson, cacheInvalidate } from '../../utils/redis.utils'; + +// ── Constants ───────────────────────────────────────────────── + +const KEY_ANALYTICS_TTL = 300; // 5 min +const PLATFORM_ANALYTICS_TTL = 300; // 5 min + +const KEY_ANALYTICS_CACHE_PREFIX = 'key:analytics:payment-assets'; +const PLATFORM_ANALYTICS_CACHE_KEY = 'platform:analytics:payment-assets'; + +// ── Cache helpers (exported for route-layer invalidation) ───── + +export function buildKeyAnalyticsCacheKey(keyId: string): string { + return `${KEY_ANALYTICS_CACHE_PREFIX}:${keyId}`; +} + +export async function invalidateKeyAnalyticsCache(keyId: string): Promise { + await cacheInvalidate( + buildKeyAnalyticsCacheKey(keyId), + PLATFORM_ANALYTICS_CACHE_KEY, + ); +} + +// ── Shared item shapes ──────────────────────────────────────── + +export interface PaymentAssetBreakdownItem { + /** Normalised asset code, e.g. 'XLM', 'USDC'. */ + paymentAsset: string; + tradeCount: number; + uniqueBuyers: number; + /** Sum of the raw `price` field (stroops) as a string to preserve precision. */ + totalPriceStroops: string; +} + +export interface KeyPaymentAssetAnalytics { + keyId: string; + generatedAt: string; + breakdown: PaymentAssetBreakdownItem[]; +} + +export interface PlatformPaymentAssetEntry extends PaymentAssetBreakdownItem { + /** Percentage share of total platform trade count, rounded to 4 dp. */ + sharePercent: number; +} + +export interface PlatformPaymentAssetDistribution { + generatedAt: string; + totalTrades: number; + breakdown: PlatformPaymentAssetEntry[]; +} + +// ── Per-key analytics ───────────────────────────────────────── + +/** + * Return trade counts, unique buyers, and total price grouped by paymentAsset + * for a single creator key. + * + * Trades with no paymentAsset (pre-migration rows) have DEFAULT 'XLM' so they + * are automatically included in the XLM bucket. + * + * Cached for 5 minutes per keyId. + */ +export async function getKeyPaymentAssetAnalytics( + keyId: string +): Promise { + const cacheKey = buildKeyAnalyticsCacheKey(keyId); + const cached = await cacheGetJson(cacheKey); + if (cached) return cached; + + // Resolve by id OR handle so the route can pass either. + const profile = await prisma.creatorProfile.findFirst({ + where: { OR: [{ id: keyId }, { handle: keyId }] }, + select: { id: true }, + }); + const resolvedId = profile?.id ?? keyId; + + // groupBy paymentAsset, counting rows and summing price. + const rows = await prisma.trade.groupBy({ + by: ['paymentAsset'], + where: { creatorId: resolvedId }, + _count: { _all: true }, + _sum: { price: false } as any, // price is String — handled via raw below + }); + + // Prisma groupBy can't SUM a String field, so we fetch the raw sums with + // findMany and aggregate in JS. The dataset per key is bounded; for very + // high volume keys this is still fast because we only pull (paymentAsset, + // price, buyer) tuples. + const trades = await prisma.trade.findMany({ + where: { creatorId: resolvedId }, + select: { paymentAsset: true, price: true, buyer: true }, + }); + + const assetMap = new Map; totalStroops: bigint }>(); + + for (const t of trades) { + const asset = t.paymentAsset; + if (!assetMap.has(asset)) { + assetMap.set(asset, { count: 0, buyers: new Set(), totalStroops: 0n }); + } + const bucket = assetMap.get(asset)!; + bucket.count++; + bucket.buyers.add(t.buyer); + try { + bucket.totalStroops += BigInt(t.price); + } catch { + // non-numeric price — skip sum contribution + } + } + + const breakdown: PaymentAssetBreakdownItem[] = Array.from(assetMap.entries()) + .map(([paymentAsset, b]) => ({ + paymentAsset, + tradeCount: b.count, + uniqueBuyers: b.buyers.size, + totalPriceStroops: b.totalStroops.toString(), + })) + .sort((a, b) => b.tradeCount - a.tradeCount); + + const result: KeyPaymentAssetAnalytics = { + keyId: resolvedId, + generatedAt: new Date().toISOString(), + breakdown, + }; + + await cacheSetJson(cacheKey, result, KEY_ANALYTICS_TTL); + return result; +} + +// ── Platform-wide distribution ──────────────────────────────── + +/** + * Return the platform-wide payment-asset distribution across all trades. + * + * Each entry includes a `sharePercent` (percentage of total trade count). + * Sorted by tradeCount descending so the dominant asset appears first. + * + * Cached for 5 minutes (single key, no per-key variation). + */ +export async function getPlatformPaymentAssetDistribution(): Promise { + const cached = await cacheGetJson(PLATFORM_ANALYTICS_CACHE_KEY); + if (cached) return cached; + + const trades = await prisma.trade.findMany({ + select: { paymentAsset: true, price: true, buyer: true }, + }); + + const assetMap = new Map; totalStroops: bigint }>(); + + for (const t of trades) { + const asset = t.paymentAsset; + if (!assetMap.has(asset)) { + assetMap.set(asset, { count: 0, buyers: new Set(), totalStroops: 0n }); + } + const bucket = assetMap.get(asset)!; + bucket.count++; + bucket.buyers.add(t.buyer); + try { + bucket.totalStroops += BigInt(t.price); + } catch { + // non-numeric price — skip + } + } + + const totalTrades = trades.length; + + const breakdown: PlatformPaymentAssetEntry[] = Array.from(assetMap.entries()) + .map(([paymentAsset, b]) => ({ + paymentAsset, + tradeCount: b.count, + uniqueBuyers: b.buyers.size, + totalPriceStroops: b.totalStroops.toString(), + sharePercent: + totalTrades > 0 + ? Math.round((b.count / totalTrades) * 100 * 10_000) / 10_000 + : 0, + })) + .sort((a, b) => b.tradeCount - a.tradeCount); + + const result: PlatformPaymentAssetDistribution = { + generatedAt: new Date().toISOString(), + totalTrades, + breakdown, + }; + + await cacheSetJson(PLATFORM_ANALYTICS_CACHE_KEY, result, PLATFORM_ANALYTICS_TTL); + return result; +} diff --git a/src/modules/keys/key-proposal-votes.service.ts b/src/modules/keys/key-proposal-votes.service.ts index ba42e707..29175b58 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) { @@ -30,7 +31,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; } /** @@ -80,9 +88,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, @@ -102,14 +117,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); } @@ -118,7 +134,20 @@ export async function castKeyProposalVote( throw new DuplicateVoteError(); } - const weight = String(balance); + // Delegated weight: sum of key balances of all active delegators. + // Fetched outside the transaction because it is read-only and does not + // need to be part of the atomic write. + 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); // TODO: submit cast_vote contract call via Stellar SDK // On-chain failure should return 502 before reaching this point. @@ -130,7 +159,10 @@ export async function castKeyProposalVote( voter: wallet, optionIndex, option: options[optionIndex], - weight, + ownWeight: ownWeightStr, + delegatedWeight: delegatedWeightStr, + totalWeight: weightStr, + delegatorCount, }, 'Submitting cast_vote contract call' ); @@ -142,21 +174,24 @@ export async function castKeyProposalVote( proposalId, voter: wallet, optionIndex, - weight: new Decimal(weight), + weight: new Decimal(weightStr), }, }), + // Use the correct GOVERNANCE_VOTE_CAST activity type (#933). prisma.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: options[optionIndex], - weight, + ownWeight: ownWeightStr, + delegatedWeight: delegatedWeightStr, + totalWeight: weightStr, + delegatorCount, }, }, }), @@ -166,6 +201,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 13c0644c..e3380eee 100644 --- a/src/modules/keys/keys.routes.ts +++ b/src/modules/keys/keys.routes.ts @@ -64,6 +64,11 @@ import { PositionNotFoundError, unfreezePosition, } from './key-freeze.service'; +import { getSunsetWatchList } from './key-sunset-watch.service'; +import { + getKeyPaymentAssetAnalytics, + getPlatformPaymentAssetDistribution, +} from './key-analytics.service'; const priceHistoryQuerySchema = z.object({ from: z.string().datetime(), @@ -184,6 +189,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 * Public key detail response includes supply milestone metadata. 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.routes.ts b/src/modules/staking/staking.routes.ts new file mode 100644 index 00000000..284854b2 --- /dev/null +++ b/src/modules/staking/staking.routes.ts @@ -0,0 +1,216 @@ +// src/modules/staking/staking.routes.ts +// HTTP route handlers for staking NFT position endpoints (#932). +// +// Endpoints: +// +// GET /api/v1/staking/nfts +// Authenticated. Returns all staking NFTs for the JWT wallet. +// Query: limit, offset, includeBurned. +// +// GET /api/v1/staking/nfts/:id +// Public. Returns full metadata for a single NFT (by DB id or tokenId). +// +// GET /api/v1/staking/nfts/:id/transfers +// Public. Returns the full ownership history for a single NFT. +// +// Route ordering: static paths (/nfts) are registered before parameterised +// paths (/:id and /:id/transfers) to avoid shadowing. + +import { Router } from 'express'; +import { z } from 'zod'; +import { + sendNotFound, + sendSuccess, + sendValidationError, + zodIssuesToDetails, +} from '../../utils/api-response.utils'; +import { + requireJwtAuth, + AuthenticatedRequest, +} from '../../middlewares/jwt-auth.middleware'; +import { logger } from '../../utils/logger.utils'; +import { + getWalletStakingNfts, + getStakingNftById, + getStakingNftTransfers, + StakingNftNotFoundError, +} from './staking.service'; + +const stakingRouter = Router(); + +// ── Shared pagination constants ────────────────────────────── + +const DEFAULT_LIMIT = 20; +const MAX_LIMIT = 100; + +// ── Zod schemas ─────────────────────────────────────────────── + +const walletNftsQuerySchema = z.object({ + limit: z.coerce + .number() + .int() + .min(1) + .max(MAX_LIMIT) + .default(DEFAULT_LIMIT), + offset: z.coerce.number().int().min(0).default(0), + includeBurned: z + .enum(['true', 'false']) + .optional() + .transform(v => v === 'true'), +}); + +const nftPaginationQuerySchema = z.object({ + limit: z.coerce + .number() + .int() + .min(1) + .max(MAX_LIMIT) + .default(DEFAULT_LIMIT), + offset: z.coerce.number().int().min(0).default(0), +}); + +// ── GET /staking/nfts ───────────────────────────────────────── + +/** + * GET /api/v1/staking/nfts + * + * Returns all staking NFT positions owned by the authenticated wallet, + * sorted by mint date descending (most recently minted first). + * + * Auth: Bearer JWT (wallet from token is the owner filter). + * + * Query parameters: + * limit — max items per page (1–100, default 20) + * offset — items to skip (default 0) + * includeBurned — include redeemed/burned NFTs (default false) + * + * Responses: + * 200 — { items: StakingNftItem[], meta: OffsetPaginationMeta } + * 400 — invalid query parameters + * 401 — missing or invalid JWT + */ +stakingRouter.get( + '/nfts', + requireJwtAuth, + async (req: AuthenticatedRequest, res, next) => { + const parsed = walletNftsQuerySchema.safeParse(req.query); + if (!parsed.success) { + sendValidationError( + res, + 'Invalid query parameters', + zodIssuesToDetails(parsed.error.issues) + ); + return; + } + + try { + const result = await getWalletStakingNfts({ + wallet: req.user!.wallet, + ...parsed.data, + }); + sendSuccess(res, result); + } catch (error) { + logger.error( + { error, wallet: req.user?.wallet }, + 'GET /staking/nfts failed' + ); + next(error); + } + } +); + +stakingRouter.all('/nfts', (_req, res) => { + res.set('Allow', 'GET').sendStatus(405); +}); + +// ── GET /staking/nfts/:id/transfers ────────────────────────── +// Must be registered BEFORE /:id so Express does not treat "transfers" +// as the :id param value and fall through to the metadata handler. + +/** + * GET /api/v1/staking/nfts/:id/transfers + * + * Returns the full ownership-transfer history for a single staking NFT, + * sorted by transfer date descending (most recent first). The initial mint + * appears as the last entry (fromAddress = null). Publicly accessible. + * + * Path parameters: + * id — the NFT's DB id or on-chain tokenId + * + * Query parameters: + * limit — max items per page (1–100, default 20) + * offset — items to skip (default 0) + * + * Responses: + * 200 — { nft: StakingNftItem, items: StakingNftTransferItem[], meta: OffsetPaginationMeta } + * 400 — invalid query parameters + * 404 — NFT not found + */ +stakingRouter.get('/:id/transfers', async (req, res, next) => { + const parsed = nftPaginationQuerySchema.safeParse(req.query); + if (!parsed.success) { + sendValidationError( + res, + 'Invalid query parameters', + zodIssuesToDetails(parsed.error.issues) + ); + return; + } + + try { + const result = await getStakingNftTransfers({ + nftIdOrTokenId: String(req.params.id), + ...parsed.data, + }); + sendSuccess(res, result); + } catch (error) { + if (error instanceof StakingNftNotFoundError) { + sendNotFound(res, 'Staking NFT'); + return; + } + logger.error( + { error, id: req.params.id }, + 'GET /staking/nfts/:id/transfers failed' + ); + next(error); + } +}); + +stakingRouter.all('/:id/transfers', (_req, res) => { + res.set('Allow', 'GET').sendStatus(405); +}); + +// ── GET /staking/nfts/:id ───────────────────────────────────── + +/** + * GET /api/v1/staking/nfts/:id + * + * Returns full metadata for a single staking NFT. Resolves by DB id or + * on-chain tokenId. Publicly accessible — no auth required. + * + * Path parameters: + * id — the NFT's DB id or on-chain tokenId + * + * Responses: + * 200 — StakingNftItem + * 404 — NFT not found + */ +stakingRouter.get('/:id', async (req, res, next) => { + try { + const nft = await getStakingNftById(String(req.params.id)); + sendSuccess(res, nft); + } catch (error) { + if (error instanceof StakingNftNotFoundError) { + sendNotFound(res, 'Staking NFT'); + return; + } + logger.error({ error, id: req.params.id }, 'GET /staking/nfts/:id failed'); + next(error); + } +}); + +stakingRouter.all('/:id', (_req, res) => { + res.set('Allow', 'GET').sendStatus(405); +}); + +export default stakingRouter; diff --git a/src/modules/staking/staking.service.ts b/src/modules/staking/staking.service.ts new file mode 100644 index 00000000..d207b949 --- /dev/null +++ b/src/modules/staking/staking.service.ts @@ -0,0 +1,278 @@ +// src/modules/staking/staking.service.ts +// Read-model service for staking NFT positions (#932). +// +// Ownership is written by the indexer (staking-indexer.service.ts) when +// StakeNFTMinted and StakeNFTTransferred contract events are processed. +// This module is read-only — it never writes to the DB directly. + +import { prisma } from '../../utils/prisma.utils'; +import { cacheGetJson, cacheSetJson } from '../../utils/redis.utils'; +import { buildOffsetPaginationMeta } from '../../utils/pagination.utils'; + +// ── Cache TTLs ─────────────────────────────────────────────── + +const WALLET_NFTS_CACHE_TTL_SECONDS = 60; +const NFT_METADATA_CACHE_TTL_SECONDS = 120; +const NFT_TRANSFERS_CACHE_TTL_SECONDS = 60; + +// ── Error classes ───────────────────────────────────────────── + +export class StakingNftNotFoundError extends Error { + constructor(nftId: string) { + super(`Staking NFT not found: ${nftId}`); + this.name = 'StakingNftNotFoundError'; + } +} + +// ── Shared item shape ───────────────────────────────────────── + +export interface StakingNftItem { + id: string; + tokenId: string; + ownerAddress: string; + keyId: string; + /** String-encoded Decimal to preserve full precision for the client. */ + stakedAmount: string; + lockExpiryLedger: number | null; + lockExpiresAt: string | null; + burned: boolean; + burnedAt: string | null; + mintLedger: number; + mintTxHash: string; + createdAt: string; + updatedAt: string; +} + +export interface StakingNftTransferItem { + id: string; + nftId: string; + fromAddress: string | null; + toAddress: string; + ledger: number; + txHash: string; + occurredAt: string; +} + +// ── Wallet NFT list ─────────────────────────────────────────── + +export interface WalletNftsQuery { + wallet: string; + /** Include burned (redeemed) NFTs in the results. Default: false. */ + includeBurned?: boolean; + limit: number; + offset: number; +} + +export interface WalletNftsResult { + items: StakingNftItem[]; + meta: ReturnType; +} + +/** + * List all staking NFTs currently owned by a wallet, sorted by creation date + * descending (most recently minted first). Non-burned positions only by + * default; pass `includeBurned: true` to include redeemed receipts. + * + * Cached per (wallet, includeBurned, limit, offset) for 60 s. + */ +export async function getWalletStakingNfts( + query: WalletNftsQuery +): Promise { + const { wallet, includeBurned = false, limit, offset } = query; + const cacheKey = `staking:nfts:wallet:${wallet}:burned_${includeBurned}:${limit}:${offset}`; + const cached = await cacheGetJson(cacheKey); + if (cached) return cached; + + const where = { + ownerAddress: wallet, + ...(includeBurned ? {} : { burned: false }), + }; + + const [rows, total] = await Promise.all([ + prisma.stakingNft.findMany({ + where, + orderBy: { createdAt: 'desc' }, + skip: offset, + take: limit, + }), + prisma.stakingNft.count({ where }), + ]); + + const result: WalletNftsResult = { + items: rows.map(mapNftRow), + meta: buildOffsetPaginationMeta({ limit, offset, total }), + }; + + await cacheSetJson(cacheKey, result, WALLET_NFTS_CACHE_TTL_SECONDS); + return result; +} + +// ── Single NFT metadata ─────────────────────────────────────── + +/** + * Fetch full metadata for a single staking NFT by its DB id or on-chain + * tokenId. Publicly accessible (no auth required at the service layer). + * + * Cached per NFT id/tokenId for 120 s. + * + * @throws {StakingNftNotFoundError} when no matching NFT exists. + */ +export async function getStakingNftById( + nftIdOrTokenId: string +): Promise { + const cacheKey = `staking:nfts:meta:${nftIdOrTokenId}`; + const cached = await cacheGetJson(cacheKey); + if (cached) return cached; + + const row = await prisma.stakingNft.findFirst({ + where: { + OR: [{ id: nftIdOrTokenId }, { tokenId: nftIdOrTokenId }], + }, + }); + + if (!row) { + throw new StakingNftNotFoundError(nftIdOrTokenId); + } + + const result = mapNftRow(row); + await cacheSetJson(cacheKey, result, NFT_METADATA_CACHE_TTL_SECONDS); + return result; +} + +// ── Transfer history ────────────────────────────────────────── + +export interface NftTransfersQuery { + nftIdOrTokenId: string; + limit: number; + offset: number; +} + +export interface NftTransfersResult { + nft: StakingNftItem; + items: StakingNftTransferItem[]; + meta: ReturnType; +} + +/** + * Return the full ownership-history for a single staking NFT, most recent + * transfer first. Resolves by DB id or on-chain tokenId. + * + * Cached per (nftIdOrTokenId, limit, offset) for 60 s. + * + * @throws {StakingNftNotFoundError} when no matching NFT exists. + */ +export async function getStakingNftTransfers( + query: NftTransfersQuery +): Promise { + const { nftIdOrTokenId, limit, offset } = query; + const cacheKey = `staking:nfts:transfers:${nftIdOrTokenId}:${limit}:${offset}`; + const cached = await cacheGetJson(cacheKey); + if (cached) return cached; + + // Resolve NFT first so we get a stable DB id for the transfer query. + const nftRow = await prisma.stakingNft.findFirst({ + where: { + OR: [{ id: nftIdOrTokenId }, { tokenId: nftIdOrTokenId }], + }, + }); + + if (!nftRow) { + throw new StakingNftNotFoundError(nftIdOrTokenId); + } + + const [transferRows, total] = await Promise.all([ + prisma.stakingNftTransfer.findMany({ + where: { nftId: nftRow.id }, + orderBy: { occurredAt: 'desc' }, + skip: offset, + take: limit, + }), + prisma.stakingNftTransfer.count({ where: { nftId: nftRow.id } }), + ]); + + const result: NftTransfersResult = { + nft: mapNftRow(nftRow), + items: transferRows.map(mapTransferRow), + meta: buildOffsetPaginationMeta({ limit, offset, total }), + }; + + await cacheSetJson(cacheKey, result, NFT_TRANSFERS_CACHE_TTL_SECONDS); + return result; +} + +// ── Cache invalidation helper ───────────────────────────────── + +/** + * Invalidate all cached entries for a wallet's NFT list. Called by the + * indexer after processing a mint or transfer event so the wallet's position + * list reflects the change within the next request. + */ +export function walletNftsCachePattern(wallet: string): string { + return `staking:nfts:wallet:${wallet}:*`; +} + +/** + * Invalidate the per-NFT metadata and transfer-history caches. Accepts both + * the DB id and the on-chain tokenId so both cache key variants are cleared. + */ +export function nftMetaCachePattern(nftId: string): string { + return `staking:nfts:meta:${nftId}`; +} + +export function nftTransfersCachePattern(nftId: string): string { + return `staking:nfts:transfers:${nftId}:*`; +} + +// ── Row mappers ─────────────────────────────────────────────── + +function mapNftRow(row: { + id: string; + tokenId: string; + ownerAddress: string; + keyId: string; + stakedAmount: { toString(): string }; + lockExpiryLedger: number | null; + lockExpiresAt: Date | null; + burned: boolean; + burnedAt: Date | null; + mintLedger: number; + mintTxHash: string; + createdAt: Date; + updatedAt: Date; +}): StakingNftItem { + return { + id: row.id, + tokenId: row.tokenId, + ownerAddress: row.ownerAddress, + keyId: row.keyId, + stakedAmount: row.stakedAmount.toString(), + lockExpiryLedger: row.lockExpiryLedger, + lockExpiresAt: row.lockExpiresAt?.toISOString() ?? null, + burned: row.burned, + burnedAt: row.burnedAt?.toISOString() ?? null, + mintLedger: row.mintLedger, + mintTxHash: row.mintTxHash, + createdAt: row.createdAt.toISOString(), + updatedAt: row.updatedAt.toISOString(), + }; +} + +function mapTransferRow(row: { + id: string; + nftId: string; + fromAddress: string | null; + toAddress: string; + ledger: number; + txHash: string; + occurredAt: Date; +}): StakingNftTransferItem { + return { + id: row.id, + nftId: row.nftId, + fromAddress: row.fromAddress, + toAddress: row.toAddress, + ledger: row.ledger, + txHash: row.txHash, + occurredAt: row.occurredAt.toISOString(), + }; +} From 747d311985452ae8dcaf49706e0fdb2bf627eb83 Mon Sep 17 00:00:00 2001 From: Harkinkunmi <158626915+TheHalalHunter@users.noreply.github.com> Date: Sat, 26 Sep 2026 20:28:44 +0100 Subject: [PATCH 2/5] fix: remove unused sendNotFound import in governance.routes.ts (ESLint #933) --- src/modules/governance/governance.routes.ts | 1 - 1 file changed, 1 deletion(-) diff --git a/src/modules/governance/governance.routes.ts b/src/modules/governance/governance.routes.ts index ee4175db..fcf43f2f 100644 --- a/src/modules/governance/governance.routes.ts +++ b/src/modules/governance/governance.routes.ts @@ -25,7 +25,6 @@ import { Router } from 'express'; import { z } from 'zod'; import { - sendNotFound, sendSuccess, sendValidationError, zodIssuesToDetails, From 5a3e8599c005e85b2e7859d0c802fc4f544360c4 Mon Sep 17 00:00:00 2001 From: Harkinkunmi <158626915+TheHalalHunter@users.noreply.github.com> Date: Sat, 26 Sep 2026 21:40:54 +0100 Subject: [PATCH 3/5] fix: resolve all CI build errors from upstream merge (issues #931-#934) --- prisma/schema/staking.prisma | 51 ++++++++++++++ .../indexer/indexer-pipeline.service.ts | 5 ++ src/modules/keys/keys.routes.ts | 68 +++++++++++++++++++ src/modules/staking/staking.service.ts | 16 +++++ 4 files changed, 140 insertions(+) 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/src/modules/indexer/indexer-pipeline.service.ts b/src/modules/indexer/indexer-pipeline.service.ts index 74423041..02c9c75a 100644 --- a/src/modules/indexer/indexer-pipeline.service.ts +++ b/src/modules/indexer/indexer-pipeline.service.ts @@ -76,6 +76,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({ diff --git a/src/modules/keys/keys.routes.ts b/src/modules/keys/keys.routes.ts index 838376c1..dcd19bca 100644 --- a/src/modules/keys/keys.routes.ts +++ b/src/modules/keys/keys.routes.ts @@ -97,6 +97,11 @@ import { matchTierForLockPeriod, calculateEffectiveWeight, } from '../staking/staking.service'; +import { getSunsetWatchList } from './key-sunset-watch.service'; +import { + getKeyPaymentAssetAnalytics, + getPlatformPaymentAssetDistribution, +} from './key-analytics.service'; const priceHistoryQuerySchema = z.object({ from: z.string().datetime(), @@ -494,6 +499,69 @@ router.get( } ); +/** + * GET /api/v1/keys/analytics/payment-assets + * Admin-only. Platform-wide payment asset distribution across all key purchases. + * Must be registered before /:keyId to avoid route shadowing. + */ +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/sunset-watch + * Admin-only. Keys approaching or past the inactivity sunset threshold. + * Must be registered before /:keyId to avoid route shadowing. + */ +const sunsetWatchQuerySchema = z.object({ + limit: z.coerce.number().int().min(1).max(100).default(20), + offset: z.coerce.number().int().min(0).default(0), +}); + +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 { + sendSuccess(res, await getSunsetWatchList({ limit: parsed.data.limit, offset: parsed.data.offset })); + } catch (error) { + logger.error({ error }, 'GET /keys/sunset-watch failed'); + next(error); + } + } +); + +/** + * GET /api/v1/keys/:keyId/analytics + * Payment-asset breakdown for a single key. Publicly accessible. + */ +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 * Public key detail response includes supply milestone metadata. 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}:*`; +} From 9704fc7d60317cd5246352f68e0c9a5f948cd7a9 Mon Sep 17 00:00:00 2001 From: Harkinkunmi <158626915+TheHalalHunter@users.noreply.github.com> Date: Sun, 27 Sep 2026 15:45:03 +0100 Subject: [PATCH 4/5] fix: repair CI and delegated vote weights --- src/modules/keys/key-analytics.service.ts | 101 +++++++++++++++++- .../keys/key-proposal-votes.service.test.ts | 38 ++++++- .../keys/key-proposal-votes.service.ts | 24 ++++- src/modules/keys/keys.routes.ts | 63 ----------- 4 files changed, 154 insertions(+), 72 deletions(-) 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 78ac69fc..1b7cbc26 100644 --- a/src/modules/keys/key-proposal-votes.service.ts +++ b/src/modules/keys/key-proposal-votes.service.ts @@ -182,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 @@ -195,7 +205,10 @@ export async function castKeyProposalVote( voter: wallet, optionIndex, option, - weight, + ownWeight: ownWeightStr, + delegatedWeight: delegatedWeightStr, + totalWeight: weightStr, + delegatorCount, }, 'Submitting cast_vote contract call' ); @@ -226,7 +239,7 @@ export async function castKeyProposalVote( totalVotingWeight: proposal.totalVotingWeight, results: proposal.results as Record, }, - weight, + weightStr, option ); @@ -258,7 +271,10 @@ export async function castKeyProposalVote( proposalId, optionIndex, option, - weight, + ownWeight: ownWeightStr, + delegatedWeight: delegatedWeightStr, + totalWeight: weightStr, + delegatorCount, }, }, }); diff --git a/src/modules/keys/keys.routes.ts b/src/modules/keys/keys.routes.ts index 4121f78b..25976636 100644 --- a/src/modules/keys/keys.routes.ts +++ b/src/modules/keys/keys.routes.ts @@ -500,69 +500,6 @@ router.get( } ); -/** - * GET /api/v1/keys/analytics/payment-assets - * Admin-only. Platform-wide payment asset distribution across all key purchases. - * Must be registered before /:keyId to avoid route shadowing. - */ -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/sunset-watch - * Admin-only. Keys approaching or past the inactivity sunset threshold. - * Must be registered before /:keyId to avoid route shadowing. - */ -const sunsetWatchQuerySchema = z.object({ - limit: z.coerce.number().int().min(1).max(100).default(20), - offset: z.coerce.number().int().min(0).default(0), -}); - -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 { - sendSuccess(res, await getSunsetWatchList({ limit: parsed.data.limit, offset: parsed.data.offset })); - } catch (error) { - logger.error({ error }, 'GET /keys/sunset-watch failed'); - next(error); - } - } -); - -/** - * GET /api/v1/keys/:keyId/analytics - * Payment-asset breakdown for a single key. Publicly accessible. - */ -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 * Public key detail response includes supply milestone metadata. From 0048b05a9d93cf71fd075927bc4ccb9bd2150d65 Mon Sep 17 00:00:00 2001 From: Harkinkunmi <158626915+TheHalalHunter@users.noreply.github.com> Date: Mon, 28 Sep 2026 12:43:38 +0100 Subject: [PATCH 5/5] fix: import key analytics route services --- src/modules/keys/keys.routes.ts | 5 +++++ 1 file changed, 5 insertions(+) diff --git a/src/modules/keys/keys.routes.ts b/src/modules/keys/keys.routes.ts index aac39d25..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(),