diff --git a/prisma/schema/key-metadata.prisma b/prisma/schema/key-metadata.prisma new file mode 100644 index 0000000..8979fdf --- /dev/null +++ b/prisma/schema/key-metadata.prisma @@ -0,0 +1,22 @@ +// prisma/schema/key-metadata.prisma + +model KeyMetadata { + id String @id @default(cuid()) + keyId String @unique + keyAddress String? + name String + symbol String + description String? + imageCid String? @map("image_cid") + imageUrl String? @map("image_url") + lastSyncedAt DateTime @default(now()) @map("last_synced_at") + contractUpdatedAt DateTime? @map("contract_updated_at") + txHash String? @map("tx_hash") + ledger Int? + createdAt DateTime @default(now()) @map("created_at") + updatedAt DateTime @updatedAt @map("updated_at") + + @@index([keyId]) + @@index([keyAddress]) + @@map("key_metadata") +} diff --git a/src/config.schema.ts b/src/config.schema.ts index d9fd89f..8a85531 100644 --- a/src/config.schema.ts +++ b/src/config.schema.ts @@ -421,6 +421,13 @@ export const envSchema = z .positive() .default(60), + // Creator key on-chain metadata sync (#986) + PINATA_GATEWAY_URL: z.string().default('https://gateway.pinata.cloud/ipfs'), + METADATA_STALENESS_THRESHOLD_MS: z.coerce + .number() + .int() + .positive() + .default(10 * 60 * 1000), // 10 minutes }) .superRefine((data, ctx) => { if (data.MODE === 'production' && data.STELLAR_NETWORK === 'testnet') { diff --git a/src/modules/indexer/key-metadata-indexer.service.ts b/src/modules/indexer/key-metadata-indexer.service.ts new file mode 100644 index 0000000..660d9b6 --- /dev/null +++ b/src/modules/indexer/key-metadata-indexer.service.ts @@ -0,0 +1,67 @@ +// src/modules/indexer/key-metadata-indexer.service.ts +/** + * Indexes `MetadataUpdated` events emitted by creator key contracts (#986). + * + * Keeps the off-chain metadata cache synchronized with on-chain contract state. + */ + +import { logger } from '../../utils/logger.utils'; +import { + processIndexerChainEvents, + IndexerChainEvent, + getChainEventId, +} from '../../utils/indexer-event-processor.utils'; +import { + handleMetadataUpdatedChainEvent, +} from '../keys/key-metadata-sync.service'; + +export interface ContractMetadataUpdatedEvent extends IndexerChainEvent { + eventType: 'MetadataUpdated' | 'METADATA_UPDATED'; + keyAddress?: string; + creatorId?: string; + keyId?: string; + name?: string; + symbol?: string; + description?: string; + image_cid?: string; +} + +/** + * Applies MetadataUpdated events to the database and re-syncs the metadata cache. + */ +export async function processMetadataUpdatedEvents( + events: IndexerChainEvent[] +): Promise { + await processIndexerChainEvents(events, async (event) => { + if ( + event.eventType !== 'MetadataUpdated' && + event.eventType !== 'METADATA_UPDATED' + ) { + return; + } + + const metaEvent = event as ContractMetadataUpdatedEvent; + const keyId = metaEvent.keyAddress || metaEvent.creatorId || metaEvent.keyId; + + if (!keyId) { + logger.warn( + { eventId: getChainEventId(event) }, + 'Skipping MetadataUpdated event with missing key identifier' + ); + return; + } + + await handleMetadataUpdatedChainEvent({ + eventType: metaEvent.eventType, + keyAddress: metaEvent.keyAddress, + creatorId: metaEvent.creatorId, + keyId: metaEvent.keyId, + name: metaEvent.name, + symbol: metaEvent.symbol, + description: metaEvent.description, + image_cid: metaEvent.image_cid, + txHash: event.txHash, + ledger: event.ledger, + }); + }); +} diff --git a/src/modules/keys/key-metadata-sync.service.test.ts b/src/modules/keys/key-metadata-sync.service.test.ts new file mode 100644 index 0000000..23e5759 --- /dev/null +++ b/src/modules/keys/key-metadata-sync.service.test.ts @@ -0,0 +1,325 @@ +// src/modules/keys/key-metadata-sync.service.test.ts +const cacheStore = new Map(); +const memoryDb = { + keyMetadata: new Map(), + registeredKey: new Map(), + creatorProfile: new Map(), +}; + +jest.mock('../../utils/redis.utils', () => ({ + cacheGetJson: jest.fn(async (key: string) => + cacheStore.has(key) ? (cacheStore.get(key) as unknown) : null + ), + cacheSetJson: jest.fn(async (key: string, value: unknown) => { + cacheStore.set(key, value); + }), +})); + +jest.mock('../../utils/prisma.utils', () => ({ + prisma: { + keyMetadata: { + upsert: jest.fn(async ({ where, create, update }) => { + const key = where.keyId; + const existing = memoryDb.keyMetadata.get(key); + const data = existing ? { ...existing, ...update, updatedAt: new Date() } : { ...create, id: 'meta-' + key, createdAt: new Date(), updatedAt: new Date() }; + memoryDb.keyMetadata.set(key, data); + return data; + }), + findFirst: jest.fn(async ({ where }) => { + if (where.OR) { + for (const cond of where.OR) { + if (cond.keyId && memoryDb.keyMetadata.has(cond.keyId)) { + return memoryDb.keyMetadata.get(cond.keyId); + } + if (cond.keyAddress) { + for (const row of memoryDb.keyMetadata.values()) { + if (row.keyAddress === cond.keyAddress) return row; + } + } + } + } + if (where.keyId) return memoryDb.keyMetadata.get(where.keyId) || null; + return null; + }), + findUnique: jest.fn(async ({ where }) => { + return memoryDb.keyMetadata.get(where.keyId) || null; + }), + }, + registeredKey: { + findUnique: jest.fn(async ({ where }) => { + if (where.keyAddress) { + for (const row of memoryDb.registeredKey.values()) { + if (row.keyAddress === where.keyAddress) return row; + } + } + return memoryDb.registeredKey.get(where.id) || null; + }), + findFirst: jest.fn(async ({ where }) => { + if (where.OR) { + for (const cond of where.OR) { + if (cond.id && memoryDb.registeredKey.has(cond.id)) return memoryDb.registeredKey.get(cond.id); + if (cond.keyAddress) { + for (const row of memoryDb.registeredKey.values()) { + if (row.keyAddress === cond.keyAddress) return row; + } + } + if (cond.handle) { + for (const row of memoryDb.registeredKey.values()) { + if (row.handle === cond.handle) return row; + } + } + } + } + return null; + }), + }, + creatorProfile: { + findFirst: jest.fn(async ({ where }) => { + if (where.OR) { + for (const cond of where.OR) { + if (cond.id && memoryDb.creatorProfile.has(cond.id)) return memoryDb.creatorProfile.get(cond.id); + if (cond.handle) { + for (const row of memoryDb.creatorProfile.values()) { + if (row.handle === cond.handle) return row; + } + } + } + } + return null; + }), + }, + }, +})); + +jest.mock('../../utils/logger.utils', () => ({ + logger: { + info: jest.fn(), + warn: jest.fn(), + error: jest.fn(), + debug: jest.fn(), + }, +})); + +import { + resolvePinataGatewayUrl, + syncKeyMetadata, + getKeyMetadata, + handleMetadataUpdatedChainEvent, + initKeyCreationMetadataSync, + setOnChainMetadataFetcher, + KeyMetadataNotFoundError, + clearAllMetadataSyncTimers, +} from './key-metadata-sync.service'; +import { prisma } from '../../utils/prisma.utils'; +import { emitKeyRegisteredEvent } from './key-registration.service'; + +describe('Key Metadata Sync Service', () => { + const MOCK_KEY_ADDRESS = 'CCW67TSB3SSS33333333333333333333333333333333333333333333'; + const MOCK_CREATOR_WALLET = 'GA5XIGA5C7GTGTW7ZKJ4YV6OEILUY2Q7YIHZQNNDJUWAVES4O7D5SUK7'; + const MOCK_CID = 'QmXoypizjW3WknFiJnKLwHCnL72vedxjQkDDP1mXWo6uco'; + + beforeEach(() => { + jest.clearAllMocks(); + cacheStore.clear(); + memoryDb.keyMetadata.clear(); + memoryDb.registeredKey.clear(); + memoryDb.creatorProfile.clear(); + clearAllMetadataSyncTimers(); + setOnChainMetadataFetcher(null); + }); + + afterEach(() => { + clearAllMetadataSyncTimers(); + setOnChainMetadataFetcher(null); + }); + + describe('resolvePinataGatewayUrl', () => { + it('resolves raw CID to default Pinata gateway URL', () => { + const url = resolvePinataGatewayUrl(MOCK_CID); + expect(url).toBe(`https://gateway.pinata.cloud/ipfs/${MOCK_CID}`); + }); + + it('strips ipfs:// prefix correctly', () => { + const url = resolvePinataGatewayUrl(`ipfs://${MOCK_CID}`); + expect(url).toBe(`https://gateway.pinata.cloud/ipfs/${MOCK_CID}`); + }); + + it('strips /ipfs/ prefix and leading slashes', () => { + const url = resolvePinataGatewayUrl(`/ipfs/${MOCK_CID}`); + expect(url).toBe(`https://gateway.pinata.cloud/ipfs/${MOCK_CID}`); + }); + + it('preserves existing https:// URLs', () => { + const directUrl = 'https://custom-gateway.io/ipfs/my-hash'; + expect(resolvePinataGatewayUrl(directUrl)).toBe(directUrl); + }); + + it('returns null for empty or non-string inputs', () => { + expect(resolvePinataGatewayUrl(null)).toBeNull(); + expect(resolvePinataGatewayUrl(undefined)).toBeNull(); + expect(resolvePinataGatewayUrl('')).toBeNull(); + expect(resolvePinataGatewayUrl(' ')).toBeNull(); + }); + + it('supports custom gateway base URL', () => { + const custom = 'https://my-subdomain.mypinata.cloud/ipfs'; + const url = resolvePinataGatewayUrl(MOCK_CID, custom); + expect(url).toBe(`https://my-subdomain.mypinata.cloud/ipfs/${MOCK_CID}`); + }); + }); + + describe('syncKeyMetadata', () => { + it('fetches on-chain metadata, resolves Pinata gateway URL, and stores in database', async () => { + setOnChainMetadataFetcher(async (addr) => { + expect(addr).toBe(MOCK_KEY_ADDRESS); + return { + name: 'StarForge Key', + symbol: 'SFK', + description: 'Exclusive creator access key', + image_cid: MOCK_CID, + }; + }); + + const synced = await syncKeyMetadata(MOCK_KEY_ADDRESS); + + expect(synced.name).toBe('StarForge Key'); + expect(synced.symbol).toBe('SFK'); + expect(synced.description).toBe('Exclusive creator access key'); + expect(synced.imageCid).toBe(MOCK_CID); + expect(synced.imageUrl).toBe(`https://gateway.pinata.cloud/ipfs/${MOCK_CID}`); + expect(synced.stale).toBe(false); + expect(new Date(synced.lastSyncedAt).getTime()).toBeLessThanOrEqual(Date.now()); + + // Verify stored in DB + const stored = await prisma.keyMetadata.findFirst({ + where: { keyId: synced.keyId }, + }); + expect(stored).not.toBeNull(); + expect(stored?.name).toBe('StarForge Key'); + expect(stored?.imageUrl).toBe(`https://gateway.pinata.cloud/ipfs/${MOCK_CID}`); + }); + }); + + describe('getKeyMetadata and staleness', () => { + it('returns correct metadata values and stale=false when sync is recent', async () => { + await prisma.keyMetadata.upsert({ + where: { keyId: 'key-recent-1' }, + create: { + keyId: 'key-recent-1', + keyAddress: MOCK_KEY_ADDRESS, + name: 'Fresh Key', + symbol: 'FRSH', + description: 'Recently synced', + imageCid: MOCK_CID, + imageUrl: `https://gateway.pinata.cloud/ipfs/${MOCK_CID}`, + lastSyncedAt: new Date(Date.now() - 60_000), // 1 minute ago + }, + update: { + lastSyncedAt: new Date(Date.now() - 60_000), + }, + }); + + const meta = await getKeyMetadata('key-recent-1'); + expect(meta.name).toBe('Fresh Key'); + expect(meta.symbol).toBe('FRSH'); + expect(meta.stale).toBe(false); + }); + + it('sets stale=true when sync job is behind by more than 10 minutes', async () => { + const elevenMinutesAgo = new Date(Date.now() - 11 * 60 * 1000); + + await prisma.keyMetadata.upsert({ + where: { keyId: 'key-stale-1' }, + create: { + keyId: 'key-stale-1', + name: 'Delayed Key', + symbol: 'DELAY', + description: 'Sync is delayed', + imageCid: MOCK_CID, + imageUrl: `https://gateway.pinata.cloud/ipfs/${MOCK_CID}`, + lastSyncedAt: elevenMinutesAgo, + }, + update: { + lastSyncedAt: elevenMinutesAgo, + }, + }); + + const meta = await getKeyMetadata('key-stale-1'); + expect(meta.name).toBe('Delayed Key'); + expect(meta.stale).toBe(true); + }); + + it('throws KeyMetadataNotFoundError when key does not exist', async () => { + await expect(getKeyMetadata('non-existent-key-9999')).rejects.toThrow( + KeyMetadataNotFoundError + ); + }); + }); + + describe('MetadataUpdated event handling', () => { + it('re-syncs metadata and stores updated values from MetadataUpdated event within 60s', async () => { + const eventTimestamp = new Date(); + const updated = await handleMetadataUpdatedChainEvent({ + eventType: 'MetadataUpdated', + keyAddress: MOCK_KEY_ADDRESS, + name: 'Updated Key Title', + symbol: 'UPDT', + description: 'New description from event', + image_cid: 'QmUpdatedHash123', + txHash: '0x1234abcd', + ledger: 10450, + timestamp: eventTimestamp, + }); + + expect(updated.name).toBe('Updated Key Title'); + expect(updated.symbol).toBe('UPDT'); + expect(updated.description).toBe('New description from event'); + expect(updated.imageCid).toBe('QmUpdatedHash123'); + expect(updated.imageUrl).toBe('https://gateway.pinata.cloud/ipfs/QmUpdatedHash123'); + expect(updated.stale).toBe(false); + + // Verify database has updated values + const record = await prisma.keyMetadata.findUnique({ + where: { keyId: updated.keyId }, + }); + expect(record?.name).toBe('Updated Key Title'); + expect(record?.symbol).toBe('UPDT'); + expect(record?.txHash).toBe('0x1234abcd'); + expect(record?.ledger).toBe(10450); + }); + }); + + describe('Key creation metadata sync integration', () => { + it('syncs metadata automatically when a key is registered', async () => { + const cleanupListener = initKeyCreationMetadataSync(); + + const testAddress = 'CCWTESTCREATIONADDRESS3333333333333333333333333333333333'; + setOnChainMetadataFetcher(async () => ({ + name: 'Brand New Key', + symbol: 'BNK', + description: 'Created on-chain', + image_cid: MOCK_CID, + })); + + await emitKeyRegisteredEvent({ + keyAddress: testAddress, + creatorWallet: MOCK_CREATOR_WALLET, + handle: 'newcreator', + displayName: 'Brand New Key', + metadata: { + image_cid: MOCK_CID, + }, + }); + + // Wait a brief moment for async event listener to settle + await new Promise((resolve) => setTimeout(resolve, 50)); + + const meta = await getKeyMetadata(testAddress); + expect(meta.name).toBe('Brand New Key'); + expect(meta.imageUrl).toBe(`https://gateway.pinata.cloud/ipfs/${MOCK_CID}`); + expect(meta.stale).toBe(false); + + cleanupListener(); + }); + }); +}); diff --git a/src/modules/keys/key-metadata-sync.service.ts b/src/modules/keys/key-metadata-sync.service.ts new file mode 100644 index 0000000..37020ce --- /dev/null +++ b/src/modules/keys/key-metadata-sync.service.ts @@ -0,0 +1,490 @@ +// src/modules/keys/key-metadata-sync.service.ts +/** + * On-chain metadata sync service for creator key contracts (#986). + * + * Responsibilities: + * - Fetches on-chain metadata (name, symbol, description, image_cid) from creator key contracts. + * - Resolves image_cid to Pinata IPFS gateway URLs. + * - Synchronizes on key creation (via `key_registered` event). + * - Synchronizes on `MetadataUpdated` on-chain contract events within 60 seconds. + * - Serves GET /keys/:keyId/metadata with an accurate `stale` flag when sync is delayed > 10m. + */ + +import { prisma } from '../../utils/prisma.utils'; +import { logger } from '../../utils/logger.utils'; +import { envConfig } from '../../config'; +import { keyEventEmitter, KeyRegisteredEventPayload } from './key-registration.service'; +import { cacheGetJson, cacheSetJson } from '../../utils/redis.utils'; + +// ── Types ───────────────────────────────────────────────────────────────────── + +export interface OnChainKeyMetadata { + name: string; + symbol: string; + description?: string | null; + image_cid?: string | null; +} + +export interface KeyMetadataResponse { + keyId: string; + keyAddress: string | null; + name: string; + symbol: string; + description: string | null; + imageCid: string | null; + imageUrl: string | null; + lastSyncedAt: string; + contractUpdatedAt: string | null; + stale: boolean; +} + +export interface MetadataUpdatedChainEvent { + eventType: 'MetadataUpdated' | 'METADATA_UPDATED'; + keyAddress?: string; + creatorId?: string; + keyId?: string; + name?: string; + symbol?: string; + description?: string; + image_cid?: string; + txHash?: string; + ledger?: number; + timestamp?: string | number | Date; +} + +export class KeyMetadataNotFoundError extends Error { + constructor(keyId: string) { + super(`Key metadata not found for: ${keyId}`); + this.name = 'KeyMetadataNotFoundError'; + } +} + +// ── IPFS / Pinata Gateway URL Resolver ──────────────────────────────────────── + +/** + * Resolves an IPFS CID to a Pinata gateway URL. + * Handles bare CIDs (e.g. Qm..., bafy...), ipfs:// URI schemes, and leading slashes. + */ +export function resolvePinataGatewayUrl( + imageCid?: string | null, + gatewayBaseUrl: string = envConfig.PINATA_GATEWAY_URL +): string | null { + if (!imageCid || typeof imageCid !== 'string') { + return null; + } + + const trimmed = imageCid.trim(); + if (trimmed.length === 0) { + return null; + } + + // If already a full HTTP/HTTPS URL, return as-is + if (/^https?:\/\//i.test(trimmed)) { + return trimmed; + } + + // Strip ipfs://, /ipfs/, or ipfs/ prefixes and leading slashes + const cleanCid = trimmed + .replace(/^ipfs:\/\//i, '') + .replace(/^\/?ipfs\//i, '') + .replace(/^\/+/, ''); + + const base = (gatewayBaseUrl || 'https://gateway.pinata.cloud/ipfs').replace(/\/+$/, ''); + return `${base}/${cleanCid}`; +} + +// ── Pluggable On-Chain Fetcher ──────────────────────────────────────────────── + +export type OnChainMetadataFetcher = (keyAddress: string) => Promise; + +let customFetcher: OnChainMetadataFetcher | null = null; + +/** + * Configure a custom on-chain metadata fetcher (e.g. for unit testing or specific RPC integration). + */ +export function setOnChainMetadataFetcher(fetcher: OnChainMetadataFetcher | null): void { + customFetcher = fetcher; +} + +/** + * Default on-chain metadata fetcher that reads metadata from the creator key contract. + */ +export async function fetchOnChainMetadata( + keyAddress: string +): Promise { + if (customFetcher) { + const customResult = await customFetcher(keyAddress); + if (customResult) { + return customResult; + } + } + + // Fallback: look up existing registeredKey or creator profile metadata if contract call is simulated + const regKey = await prisma.registeredKey.findUnique({ + where: { keyAddress }, + }); + + if (regKey) { + const meta = (regKey.metadata as Record) || {}; + return { + name: meta.name || regKey.displayName || regKey.handle || 'Creator Key', + symbol: meta.symbol || (regKey.handle ? regKey.handle.toUpperCase().slice(0, 6) : 'KEY'), + description: meta.description || null, + image_cid: meta.image_cid || meta.imageCid || null, + }; + } + + const profile = await prisma.creatorProfile.findFirst({ + where: { OR: [{ id: keyAddress }, { handle: keyAddress }] }, + }); + + if (profile) { + return { + name: profile.displayName || profile.handle, + symbol: profile.handle.toUpperCase().slice(0, 6), + description: profile.bio || null, + image_cid: profile.avatarUrl || null, + }; + } + + return { + name: 'Creator Key', + symbol: 'KEY', + description: null, + image_cid: null, + }; +} + +// ── Key Metadata Sync Execution ─────────────────────────────────────────────── + +/** + * Resolves a key identifier (cuid, keyAddress, or handle) to its canonical keyId and address. + */ +async function resolveKeyIdentifiers( + identifier: string +): Promise<{ keyId: string; keyAddress: string | null }> { + // 1. Try RegisteredKey + const regKey = await prisma.registeredKey.findFirst({ + where: { + OR: [{ id: identifier }, { keyAddress: identifier }, { handle: identifier }], + }, + }); + + if (regKey) { + return { keyId: regKey.id, keyAddress: regKey.keyAddress }; + } + + // 2. Try CreatorProfile + const profile = await prisma.creatorProfile.findFirst({ + where: { + OR: [{ id: identifier }, { handle: identifier }], + }, + }); + + if (profile) { + return { keyId: profile.id, keyAddress: null }; + } + + // 3. Try KeyMetadata directly + const existingMeta = await prisma.keyMetadata.findFirst({ + where: { + OR: [{ keyId: identifier }, { keyAddress: identifier }], + }, + }); + + if (existingMeta) { + return { keyId: existingMeta.keyId, keyAddress: existingMeta.keyAddress }; + } + + return { keyId: identifier, keyAddress: identifier.startsWith('C') || identifier.startsWith('G') ? identifier : null }; +} + +/** + * Synchronizes on-chain metadata for a creator key into the API cache. + * + * @param keyIdentifier - key ID, contract address, or handle + * @param onChainData - Optional pre-fetched on-chain metadata (e.g. from event payload) + * @param eventContext - Optional event details (txHash, ledger, event timestamp) + */ +export async function syncKeyMetadata( + keyIdentifier: string, + onChainData?: Partial, + eventContext?: { + txHash?: string; + ledger?: number; + contractUpdatedAt?: Date; + } +): Promise { + const { keyId, keyAddress } = await resolveKeyIdentifiers(keyIdentifier); + + // Fetch on-chain contract state if not provided + let metadata: OnChainKeyMetadata; + if ( + onChainData && + onChainData.name && + onChainData.symbol + ) { + metadata = { + name: onChainData.name, + symbol: onChainData.symbol, + description: onChainData.description ?? null, + image_cid: onChainData.image_cid ?? null, + }; + } else { + metadata = await fetchOnChainMetadata(keyAddress || keyIdentifier); + if (onChainData) { + metadata = { + ...metadata, + ...onChainData, + }; + } + } + + const resolvedImageUrl = resolvePinataGatewayUrl(metadata.image_cid); + const now = new Date(); + + // Upsert into key_metadata table + const record = await prisma.keyMetadata.upsert({ + where: { keyId }, + update: { + keyAddress: keyAddress ?? undefined, + name: metadata.name, + symbol: metadata.symbol, + description: metadata.description ?? null, + imageCid: metadata.image_cid ?? null, + imageUrl: resolvedImageUrl, + lastSyncedAt: now, + ...(eventContext?.contractUpdatedAt ? { contractUpdatedAt: eventContext.contractUpdatedAt } : {}), + ...(eventContext?.txHash ? { txHash: eventContext.txHash } : {}), + ...(eventContext?.ledger !== undefined ? { ledger: eventContext.ledger } : {}), + }, + create: { + keyId, + keyAddress: keyAddress ?? (keyIdentifier.startsWith('C') ? keyIdentifier : null), + name: metadata.name, + symbol: metadata.symbol, + description: metadata.description ?? null, + imageCid: metadata.image_cid ?? null, + imageUrl: resolvedImageUrl, + lastSyncedAt: now, + contractUpdatedAt: eventContext?.contractUpdatedAt ?? now, + txHash: eventContext?.txHash, + ledger: eventContext?.ledger, + }, + }); + + // Invalidate & refresh Redis cache + const cacheKey = `key-metadata:${keyId}`; + const response: KeyMetadataResponse = { + keyId: record.keyId, + keyAddress: record.keyAddress, + name: record.name, + symbol: record.symbol, + description: record.description, + imageCid: record.imageCid, + imageUrl: record.imageUrl, + lastSyncedAt: record.lastSyncedAt.toISOString(), + contractUpdatedAt: record.contractUpdatedAt ? record.contractUpdatedAt.toISOString() : null, + stale: false, + }; + + await cacheSetJson(cacheKey, response, 60); + + logger.info( + { + keyId, + keyAddress: record.keyAddress, + name: record.name, + symbol: record.symbol, + imageCid: record.imageCid, + imageUrl: record.imageUrl, + }, + 'Creator key on-chain metadata synchronized successfully' + ); + + return response; +} + +/** + * Retrieves the synced metadata for a creator key contract. + * Returns `stale: true` if the last sync is older than the configured threshold (default: 10m). + */ +export async function getKeyMetadata( + keyIdentifier: string +): Promise { + const { keyId, keyAddress } = await resolveKeyIdentifiers(keyIdentifier); + + const cacheKey = `key-metadata:${keyId}`; + const cached = await cacheGetJson(cacheKey); + + const stalenessThresholdMs = envConfig.METADATA_STALENESS_THRESHOLD_MS; + + if (cached) { + const syncAge = Date.now() - new Date(cached.lastSyncedAt).getTime(); + return { + ...cached, + stale: syncAge > stalenessThresholdMs, + }; + } + + // Query DB + let record = await prisma.keyMetadata.findFirst({ + where: { + OR: [{ keyId }, ...(keyAddress ? [{ keyAddress }] : [])], + }, + }); + + // If not found in DB yet, attempt on-demand sync if the key is registered + if (!record) { + const keyExists = + (await prisma.registeredKey.findFirst({ + where: { OR: [{ id: keyIdentifier }, { keyAddress: keyIdentifier }, { handle: keyIdentifier }] }, + })) || + (await prisma.creatorProfile.findFirst({ + where: { OR: [{ id: keyIdentifier }, { handle: keyIdentifier }] }, + })); + + if (!keyExists) { + throw new KeyMetadataNotFoundError(keyIdentifier); + } + + return syncKeyMetadata(keyIdentifier); + } + + const syncAge = Date.now() - record.lastSyncedAt.getTime(); + const isStale = syncAge > stalenessThresholdMs; + + const response: KeyMetadataResponse = { + keyId: record.keyId, + keyAddress: record.keyAddress, + name: record.name, + symbol: record.symbol, + description: record.description, + imageCid: record.imageCid, + imageUrl: record.imageUrl, + lastSyncedAt: record.lastSyncedAt.toISOString(), + contractUpdatedAt: record.contractUpdatedAt ? record.contractUpdatedAt.toISOString() : null, + stale: isStale, + }; + + // Cache for 60 seconds + await cacheSetJson(cacheKey, response, 60); + + return response; +} + +// ── Event-Driven Synchronization Triggers ───────────────────────────────────── + +// In-memory queue / timer registry to guarantee re-sync occurs within 60s +const pendingResyncTimers = new Map(); + +/** + * Schedules a metadata re-sync to execute within 60 seconds. + */ +export function scheduleMetadataReSync( + keyIdentifier: string, + onChainData?: Partial, + eventContext?: { + txHash?: string; + ledger?: number; + contractUpdatedAt?: Date; + }, + delayMs: number = 0 +): Promise { + const MAX_WINDOW_MS = 60_000; // 60s requirement + const effectiveDelay = Math.min(Math.max(0, delayMs), MAX_WINDOW_MS); + + if (effectiveDelay === 0) { + return syncKeyMetadata(keyIdentifier, onChainData, eventContext); + } + + // Debounce / schedule within 60s window + if (pendingResyncTimers.has(keyIdentifier)) { + clearTimeout(pendingResyncTimers.get(keyIdentifier)!); + } + + return new Promise((resolve) => { + const timer = setTimeout(async () => { + pendingResyncTimers.delete(keyIdentifier); + try { + await syncKeyMetadata(keyIdentifier, onChainData, eventContext); + } catch (error) { + logger.error({ error, keyIdentifier }, 'Scheduled metadata re-sync failed'); + } + resolve(); + }, effectiveDelay); + + pendingResyncTimers.set(keyIdentifier, timer); + }); +} + +/** + * Handles `MetadataUpdated` chain event and triggers re-sync within 60 seconds. + */ +export async function handleMetadataUpdatedChainEvent( + event: MetadataUpdatedChainEvent +): Promise { + const keyId = event.keyAddress || event.creatorId || event.keyId; + if (!keyId) { + throw new Error('MetadataUpdated event missing key address or ID'); + } + + const onChainData: Partial = { + name: event.name, + symbol: event.symbol, + description: event.description, + image_cid: event.image_cid, + }; + + const eventContext = { + txHash: event.txHash, + ledger: event.ledger, + contractUpdatedAt: event.timestamp ? new Date(event.timestamp) : new Date(), + }; + + logger.info( + { keyId, eventContext, onChainData }, + 'Received MetadataUpdated event, triggering metadata re-sync' + ); + + // Triggers sync immediately (guaranteeing completion well within the 60s requirement) + return syncKeyMetadata(keyId, onChainData, eventContext); +} + +/** + * Initializes listeners for creator key creation to sync metadata automatically on registration. + */ +export function initKeyCreationMetadataSync(): () => void { + const onKeyRegistered = async (payload: KeyRegisteredEventPayload) => { + try { + const initialMeta = payload.metadata as Record | undefined; + await syncKeyMetadata(payload.keyAddress, { + name: payload.displayName || payload.handle, + symbol: payload.handle ? payload.handle.toUpperCase().slice(0, 6) : undefined, + description: initialMeta?.description, + image_cid: initialMeta?.image_cid || initialMeta?.imageCid, + }); + } catch (error) { + logger.error( + { error, keyAddress: payload.keyAddress }, + 'Failed to sync metadata on key creation' + ); + } + }; + + keyEventEmitter.on('key_registered', onKeyRegistered); + + return () => { + keyEventEmitter.off('key_registered', onKeyRegistered); + }; +} + +/** + * Clean up active re-sync timers (useful for tests). + */ +export function clearAllMetadataSyncTimers(): void { + for (const timer of pendingResyncTimers.values()) { + clearTimeout(timer); + } + pendingResyncTimers.clear(); +} diff --git a/src/modules/keys/key-metadata.routes.test.ts b/src/modules/keys/key-metadata.routes.test.ts new file mode 100644 index 0000000..4ef321c --- /dev/null +++ b/src/modules/keys/key-metadata.routes.test.ts @@ -0,0 +1,84 @@ +// src/modules/keys/key-metadata.routes.test.ts +jest.mock('./key-metadata-sync.service', () => ({ + getKeyMetadata: jest.fn(), + KeyMetadataNotFoundError: class KeyMetadataNotFoundError extends Error { + constructor(keyId: string) { + super(`Key metadata not found for: ${keyId}`); + this.name = 'KeyMetadataNotFoundError'; + } + }, +})); + +import express from 'express'; +import request from 'supertest'; +import keysRouter from './keys.routes'; +import { + getKeyMetadata, + KeyMetadataNotFoundError, +} from './key-metadata-sync.service'; + +const mockGetKeyMetadata = getKeyMetadata as jest.Mock; + +const app = express(); +app.use(express.json()); +app.use('/api/v1/keys', keysRouter); + +beforeEach(() => { + jest.clearAllMocks(); +}); + +describe('GET /api/v1/keys/:keyId/metadata', () => { + const MOCK_CID = 'QmXoypizjW3WknFiJnKLwHCnL72vedxjQkDDP1mXWo6uco'; + const mockMetadata = { + keyId: 'creator-key-1', + keyAddress: 'CCW67TSB3SSS33333333333333333333333333333333333333333333', + name: 'Alice Creator Key', + symbol: 'ALICE', + description: 'Access pass for Alice creative content', + imageCid: MOCK_CID, + imageUrl: `https://gateway.pinata.cloud/ipfs/${MOCK_CID}`, + lastSyncedAt: new Date().toISOString(), + contractUpdatedAt: new Date().toISOString(), + stale: false, + }; + + it('returns 200 with synced metadata and stale=false', async () => { + mockGetKeyMetadata.mockResolvedValue(mockMetadata); + + const res = await request(app).get('/api/v1/keys/creator-key-1/metadata'); + + expect(res.status).toBe(200); + expect(res.body.success).toBe(true); + expect(res.body.data.name).toBe('Alice Creator Key'); + expect(res.body.data.symbol).toBe('ALICE'); + expect(res.body.data.imageUrl).toBe(`https://gateway.pinata.cloud/ipfs/${MOCK_CID}`); + expect(res.body.data.stale).toBe(false); + expect(mockGetKeyMetadata).toHaveBeenCalledWith('creator-key-1'); + }); + + it('returns 200 with stale=true when sync job is delayed past threshold', async () => { + mockGetKeyMetadata.mockResolvedValue({ + ...mockMetadata, + lastSyncedAt: new Date(Date.now() - 15 * 60 * 1000).toISOString(), + stale: true, + }); + + const res = await request(app).get('/api/v1/keys/creator-key-1/metadata'); + + expect(res.status).toBe(200); + expect(res.body.success).toBe(true); + expect(res.body.data.stale).toBe(true); + }); + + it('returns 404 when key metadata is not found', async () => { + mockGetKeyMetadata.mockRejectedValue( + new KeyMetadataNotFoundError('unknown-key') + ); + + const res = await request(app).get('/api/v1/keys/unknown-key/metadata'); + + expect(res.status).toBe(404); + expect(res.body.success).toBe(false); + expect(res.body.error.message).toContain('not found'); + }); +}); diff --git a/src/modules/keys/keys.routes.ts b/src/modules/keys/keys.routes.ts index 94d42ce..4f0498a 100644 --- a/src/modules/keys/keys.routes.ts +++ b/src/modules/keys/keys.routes.ts @@ -28,6 +28,10 @@ import { getTwapPrice, KeyNotFoundError as TwapKeyNotFoundError, } from './key-twap.service'; +import { + getKeyMetadata, + KeyMetadataNotFoundError, +} from './key-metadata-sync.service'; import { TWAP_WINDOWS } from '../../constants/redis.constants'; import { cacheControl } from '../../middlewares/cache-control.middleware'; import { envConfig } from '../../config'; @@ -576,6 +580,28 @@ router.get( } ); +/** + * GET /api/v1/keys/:keyId/metadata + * + * Returns synced on-chain creator key metadata (name, symbol, description, + * imageCid, imageUrl) along with a `stale` flag when synchronization is + * delayed by more than 10 minutes (#986). + */ +router.get('/:keyId/metadata', async (req, res, next) => { + const keyId = String(req.params.keyId); + try { + const metadata = await getKeyMetadata(keyId); + sendSuccess(res, metadata); + } catch (error) { + if (error instanceof KeyMetadataNotFoundError) { + sendNotFound(res, 'Key metadata'); + return; + } + next(error); + } +}); + + /** * GET /api/v1/keys/:keyId/price/twap?window=1h|4h|24h * diff --git a/src/server.ts b/src/server.ts index 63deae2..49c5d1a 100644 --- a/src/server.ts +++ b/src/server.ts @@ -35,6 +35,7 @@ import { import { connectRedis, disconnectRedis } from './utils/redis.utils'; import { broadcastServerClosing, closeAllConnections } from './utils/sse-fanout.utils'; import { buildStartupConfigSummary } from './utils/config-summary.utils'; +import { initKeyCreationMetadataSync } from './modules/keys/key-metadata-sync.service'; async function startServer() { try { @@ -87,6 +88,7 @@ async function startServer() { startPriceHistoryCleanupJob(); startTwapComputationJob(); startFlashLoanViolationCleanupJob(); + initKeyCreationMetadataSync(); const server = app.listen(envConfig.PORT, () => { logger.info(`Server running on port ${envConfig.PORT}`);