diff --git a/backend/.env.example b/backend/.env.example index 95a5582..3353ae1 100644 --- a/backend/.env.example +++ b/backend/.env.example @@ -15,10 +15,14 @@ DATABASE_CONNECT_TIMEOUT=10 PORT=4000 NODE_ENV=development -# Stellar / Indexer -INDEXER_INTERVAL_MS=5000 +# Stellar / Indexer (interval must stay ≤ 2000ms for PaymentStreamed freshness) +INDEXER_INTERVAL_MS=1000 STELLAR_NETWORK=testnet STELLAR_HORIZON_URL=https://horizon-testnet.stellar.org +SOROBAN_RPC_URL=https://soroban-testnet.stellar.org +# Comma-separated contract IDs used to filter getEvents (required in production) +ESCROW_CONTRACT_ID= +DISPUTE_CONTRACT_ID= # Expert notifications (Discord / Telegram) — optional; failures never block the API # Discord channel incoming webhook URL diff --git a/backend/prisma/migrations/20260822120000_indexer_cursor/migration.sql b/backend/prisma/migrations/20260822120000_indexer_cursor/migration.sql new file mode 100644 index 0000000..39dde92 --- /dev/null +++ b/backend/prisma/migrations/20260822120000_indexer_cursor/migration.sql @@ -0,0 +1,2 @@ +-- AlterTable +ALTER TABLE "IndexerState" ADD COLUMN "lastCursor" TEXT; diff --git a/backend/prisma/schema.prisma b/backend/prisma/schema.prisma index 3a37fb2..55f8689 100644 --- a/backend/prisma/schema.prisma +++ b/backend/prisma/schema.prisma @@ -179,7 +179,7 @@ model Transaction { model EventLog { id String @id @default(cuid()) txHash String @unique - /// e.g. "SESSION_BOOKED", "SESSION_COMPLETED", "PAYMENT_RELEASED" + /// e.g. "SESSION_BOOKED", "PAYMENT_STREAMED", "PAYMENT_RELEASED" eventType String /// JSON-encoded event payload payload String @@ -198,6 +198,8 @@ model EventLog { model IndexerState { key String @id lastLedger Int @default(0) + /// Opaque getEvents cursor for pagination (mutually exclusive with startLedger) + lastCursor String? updatedAt DateTime @updatedAt } diff --git a/backend/src/eventListener.ts b/backend/src/eventListener.ts index 3e283f1..b309d64 100644 --- a/backend/src/eventListener.ts +++ b/backend/src/eventListener.ts @@ -1,4 +1,4 @@ -import { PrismaClient } from "@prisma/client"; +import { Prisma, PrismaClient } from "@prisma/client"; import { getNotificationService } from "./notificationService"; import { publishSessionStatus, @@ -9,6 +9,7 @@ export type StellarEventType = | "SESSION_BOOKED" | "SESSION_COMPLETED" | "PAYMENT_RELEASED" + | "PAYMENT_STREAMED" | "EXPERT_REGISTERED" | "SESSION_PAUSED" | "SESSION_REFUNDED"; @@ -18,6 +19,7 @@ const SESSION_SCOPED_EVENTS: ReadonlySet = new Set([ "SESSION_BOOKED", "SESSION_COMPLETED", "PAYMENT_RELEASED", + "PAYMENT_STREAMED", "SESSION_PAUSED", "SESSION_REFUNDED", ]); @@ -50,16 +52,30 @@ export async function ingestEvent( return { created: false, id: existing.id }; } - const created = await prisma.eventLog.create({ - data: { - txHash: event.txHash, - eventType: event.eventType, - payload: JSON.stringify(event.payload), - processed: false, - }, - }); - - return { created: true, id: created.id }; + try { + const created = await prisma.eventLog.create({ + data: { + txHash: event.txHash, + eventType: event.eventType, + payload: JSON.stringify(event.payload), + processed: false, + }, + }); + return { created: true, id: created.id }; + } catch (err) { + if ( + err instanceof Prisma.PrismaClientKnownRequestError && + err.code === "P2002" + ) { + const raced = await prisma.eventLog.findUnique({ + where: { txHash: event.txHash }, + }); + if (raced) { + return { created: false, id: raced.id }; + } + } + throw err; + } } /** @@ -131,6 +147,9 @@ async function handleEvent( case "PAYMENT_RELEASED": await handlePaymentReleased(prisma, payload, txHash); break; + case "PAYMENT_STREAMED": + await handlePaymentStreamed(prisma, payload, txHash); + break; case "SESSION_PAUSED": await handleSessionPaused(prisma, payload); break; @@ -249,9 +268,8 @@ async function handleSessionBooked( }); } - const hash = txHash ?? (payload["txHash"] as string); + const hash = transactionHash(payload, txHash); if (hash) { - const ledgerTime = payload["ledgerTime"] ? new Date(payload["ledgerTime"] as string) : new Date(); await prisma.transaction.upsert({ where: { txHash: hash }, create: { @@ -259,7 +277,7 @@ async function handleSessionBooked( sessionId, amount: escrowAmount, type: "ESCROW_FUNDED", - ledgerTime, + ledgerTime: ledgerTimeFromPayload(payload), }, update: {}, }); @@ -279,11 +297,10 @@ async function handleSessionCompleted( data: { status: "COMPLETED", endTime: new Date() }, }); - const hash = txHash ?? (payload["txHash"] as string); + const hash = transactionHash(payload, txHash); const amountStr = payload["amount"] ?? payload["escrowAmount"]; if (hash && amountStr !== undefined) { const amount = BigInt(String(amountStr)); - const ledgerTime = payload["ledgerTime"] ? new Date(payload["ledgerTime"] as string) : new Date(); await prisma.transaction.upsert({ where: { txHash: hash }, create: { @@ -291,7 +308,7 @@ async function handleSessionCompleted( sessionId, amount, type: "PAYMENT_RELEASED", - ledgerTime, + ledgerTime: ledgerTimeFromPayload(payload), }, update: {}, }); @@ -317,9 +334,8 @@ async function handlePaymentReleased( }); } - const hash = txHash ?? (payload["txHash"] as string); + const hash = transactionHash(payload, txHash); if (hash && sessionId) { - const ledgerTime = payload["ledgerTime"] ? new Date(payload["ledgerTime"] as string) : new Date(); await prisma.transaction.upsert({ where: { txHash: hash }, create: { @@ -327,13 +343,58 @@ async function handlePaymentReleased( sessionId, amount, type: "PAYMENT_RELEASED", - ledgerTime, + ledgerTime: ledgerTimeFromPayload(payload), }, update: {}, }); } } +/** + * Incremental expert earnings during an active session (on-chain PaymentStreamed). + * Records a PAYMENT_RELEASED transaction without completing the session. + */ +async function handlePaymentStreamed( + prisma: PrismaClient, + payload: Record, + txHash?: string +): Promise { + if (!payload["sessionId"]) throw new Error("PAYMENT_STREAMED: missing sessionId"); + if (payload["amount"] === undefined || payload["amount"] === null) { + throw new Error("PAYMENT_STREAMED: missing amount"); + } + + const sessionId = String(payload["sessionId"]); + const amount = BigInt(String(payload["amount"])); + + const session = await prisma.session.findUnique({ + where: { sessionId }, + }); + if (!session) { + throw new Error(`PAYMENT_STREAMED: session not found: ${sessionId}`); + } + + const hash = transactionHash(payload, txHash); + if (!hash) { + throw new Error("PAYMENT_STREAMED: missing txHash"); + } + + await prisma.transaction.upsert({ + where: { txHash: hash }, + create: { + txHash: hash, + sessionId, + amount, + type: "PAYMENT_RELEASED", + ledgerTime: ledgerTimeFromPayload(payload), + }, + update: { + amount, + ledgerTime: ledgerTimeFromPayload(payload), + }, + }); +} + async function handleSessionPaused( prisma: PrismaClient, payload: Record @@ -360,11 +421,10 @@ async function handleSessionRefunded( data: { status: "REFUNDED" }, }); - const hash = txHash ?? (payload["txHash"] as string); + const hash = transactionHash(payload, txHash); const amountStr = payload["amount"] ?? payload["escrowAmount"]; if (hash && amountStr !== undefined) { const amount = BigInt(String(amountStr)); - const ledgerTime = payload["ledgerTime"] ? new Date(payload["ledgerTime"] as string) : new Date(); await prisma.transaction.upsert({ where: { txHash: hash }, create: { @@ -372,13 +432,35 @@ async function handleSessionRefunded( sessionId, amount, type: "REFUND_ISSUED", - ledgerTime, + ledgerTime: ledgerTimeFromPayload(payload), }, update: {}, }); } } +function transactionHash( + payload: Record, + txHash?: string +): string | undefined { + if (typeof payload["txHash"] === "string" && payload["txHash"]) { + return payload["txHash"]; + } + if (typeof txHash === "string" && txHash) { + return txHash; + } + return undefined; +} + +function ledgerTimeFromPayload(payload: Record): Date { + const raw = payload["ledgerClosedAt"] ?? payload["ledgerTime"]; + if (typeof raw === "string" && raw) { + const parsed = new Date(raw); + if (!Number.isNaN(parsed.getTime())) return parsed; + } + return new Date(); +} + /** * If the indexer payload includes a sessionId, fan out to the session room. * Producers without a sessionId (e.g. early SESSION_BOOKED stubs) are skipped. @@ -395,6 +477,8 @@ function maybeBroadcastSessionStatus( let wsType: SessionStatusEventType = eventType as SessionStatusEventType; if (eventType === "SESSION_REFUNDED") { wsType = "SESSION_ENDED"; + } else if (eventType === "PAYMENT_STREAMED") { + wsType = "PAYMENT_RELEASED"; } publishSessionStatus( diff --git a/backend/src/indexer.ts b/backend/src/indexer.ts index 98c348d..8c5576c 100644 --- a/backend/src/indexer.ts +++ b/backend/src/indexer.ts @@ -2,27 +2,37 @@ import { prisma } from "./prisma"; import { processEvents } from "./eventListener"; import { SorobanIndexerService } from "./sorobanIndexer"; -const INTERVAL_MS = Number(process.env.INDEXER_INTERVAL_MS ?? 5000); +const DEFAULT_INTERVAL_MS = 1000; +const MIN_INTERVAL_MS = 100; + +function resolveIntervalMs(raw: string | undefined): number { + const parsed = Number(raw ?? DEFAULT_INTERVAL_MS); + if (!Number.isFinite(parsed) || parsed < MIN_INTERVAL_MS) { + return DEFAULT_INTERVAL_MS; + } + return Math.floor(parsed); +} + +const INTERVAL_MS = resolveIntervalMs(process.env.INDEXER_INTERVAL_MS); const sleep = (ms: number) => new Promise((resolve) => setTimeout(resolve, ms)); const sorobanIndexer = new SorobanIndexerService(prisma); async function tick(): Promise { - // Poll Stellar / Soroban RPC endpoint for contract events const sorobanRes = await sorobanIndexer.pollOnce(); - - // Process any pending / unhandled event logs + // Catch events ingested if processEvents inside pollOnce was interrupted. const result = await processEvents(prisma); + const processed = sorobanRes.processedCount + result.processed; if ( sorobanRes.eventsFetched > 0 || - result.processed > 0 || + processed > 0 || result.skipped > 0 || result.errors.length > 0 ) { console.log( - `[indexer] sorobanEvents=${sorobanRes.eventsFetched} processed=${result.processed} skipped=${result.skipped} errors=${result.errors.length} latestLedger=${sorobanRes.latestLedger}` + `[indexer] sorobanEvents=${sorobanRes.eventsFetched} processed=${processed} skipped=${result.skipped} errors=${result.errors.length} latestLedger=${sorobanRes.latestLedger}` ); for (const error of result.errors) { @@ -32,11 +42,8 @@ async function tick(): Promise { } async function main(): Promise { - console.log( - `SkillSphere indexer running (interval=${INTERVAL_MS}ms)` - ); + console.log(`SkillSphere indexer running (interval=${INTERVAL_MS}ms)`); - // Verify DB connectivity before entering the poll loop await prisma.$connect(); console.log("[indexer] connected to database"); diff --git a/backend/src/sorobanIndexer.ts b/backend/src/sorobanIndexer.ts index 77c17a0..caad894 100644 --- a/backend/src/sorobanIndexer.ts +++ b/backend/src/sorobanIndexer.ts @@ -8,10 +8,34 @@ import { } from "./eventListener"; const INDEXER_KEY = "stellar_soroban_indexer"; +const GET_EVENTS_LIMIT = 100; +const MAX_PAGES_PER_TICK = 10; +const MAX_CONTRACT_IDS_PER_FILTER = 5; +const MAX_FILTERS = 5; +const LOOKBACK_LEDGERS = 10; + +const EVENT_TYPE_BY_TOPIC: Record = { + paymentstreamed: "PAYMENT_STREAMED", + streamedpayment: "PAYMENT_STREAMED", + refundsession: "SESSION_REFUNDED", + sessionrefunded: "SESSION_REFUNDED", + refund: "SESSION_REFUNDED", + fundsession: "SESSION_BOOKED", + sessionbooked: "SESSION_BOOKED", + booked: "SESSION_BOOKED", + pausesession: "SESSION_PAUSED", + sessionpaused: "SESSION_PAUSED", + pause: "SESSION_PAUSED", + completesession: "SESSION_COMPLETED", + sessioncompleted: "SESSION_COMPLETED", + paymentreleased: "SESSION_COMPLETED", + complete: "SESSION_COMPLETED", + expertregistered: "EXPERT_REGISTERED", + registerexpert: "EXPERT_REGISTERED", +}; /** - * Utility function to convert scVal XDR (either base64 string or xdr.ScVal) - * to native JS primitives or objects. + * Convert scVal XDR (base64 string or xdr.ScVal) to native JS values. */ export function parseScVal(val: unknown): unknown { if (val === null || val === undefined) return val; @@ -43,6 +67,16 @@ export function toPlainObject(val: unknown): unknown { if (typeof val === "bigint") { return val.toString(); } + if (typeof val === "symbol") { + return val.description ?? val.toString(); + } + if (val instanceof Uint8Array) { + try { + return new TextDecoder().decode(val); + } catch { + return Buffer.from(val).toString("base64"); + } + } if (val instanceof Map) { const obj: Record = {}; for (const [k, v] of val.entries()) { @@ -63,9 +97,63 @@ export function toPlainObject(val: unknown): unknown { return val; } +function topicToName(topic: unknown): string { + if (typeof topic === "symbol") { + return topic.description ?? ""; + } + if (typeof topic === "string") return topic; + if (typeof topic === "number" || typeof topic === "bigint") { + return String(topic); + } + if (topic instanceof Uint8Array) { + return new TextDecoder().decode(topic); + } + if (topic == null) return ""; + const plain = toPlainObject(topic); + return typeof plain === "string" ? plain : ""; +} + +function normalizeEventName(name: string): string { + return name.toLowerCase().replace(/[^a-z0-9]/g, ""); +} + +/** + * Map a decoded Soroban topic name onto the off-chain event union. + * Names must match exactly after normalization (no substring matching). + */ +export function classifyEventType(topicName: string): StellarEventType | null { + const n = normalizeEventName(topicName); + if (!n) return null; + return EVENT_TYPE_BY_TOPIC[n] ?? null; +} + +function assignAliasedFields(payload: Record): void { + const aliases: Array<[string, string]> = [ + ["session_id", "sessionId"], + ["seeker_address", "seekerAddress"], + ["expert_address", "expertAddress"], + ["expert_id", "expertId"], + ["escrow_amount", "amount"], + ["wallet_address", "walletAddress"], + ]; + for (const [from, to] of aliases) { + if (payload[from] !== undefined && payload[to] === undefined) { + payload[to] = payload[from]; + } + } +} + +function asScalar(value: unknown): string | number | boolean | null { + if (value === null || value === undefined) return null; + if (typeof value === "string" || typeof value === "number" || typeof value === "boolean") { + return value; + } + if (typeof value === "bigint") return value.toString(); + return null; +} + /** - * Safely decode XDR topics and value payloads into a typed StellarEvent structure. - * Catches any XDR parsing or mapping errors to ensure service resilience. + * Decode XDR topics and value payloads into a typed StellarEvent structure. */ export function decodeEventPayload( topicRaw: unknown[], @@ -76,65 +164,43 @@ export function decodeEventPayload( parseScVal(t) ); const decodedValue = parseScVal(valueRaw); - - const topicName = String(decodedTopics[0] ?? "").toLowerCase(); - - let eventType: StellarEventType | null = null; - if ( - topicName.includes("refund_session") || - topicName.includes("refund") || - topicName.includes("session_refunded") - ) { - eventType = "SESSION_REFUNDED"; - } else if ( - topicName.includes("fund_session") || - topicName.includes("fund") || - topicName.includes("session_booked") || - topicName.includes("booked") - ) { - eventType = "SESSION_BOOKED"; - } else if ( - topicName.includes("pause_session") || - topicName.includes("pause") || - topicName.includes("session_paused") - ) { - eventType = "SESSION_PAUSED"; - } else if ( - topicName.includes("complete_session") || - topicName.includes("complete") || - topicName.includes("payment_released") || - topicName.includes("payment") - ) { - eventType = "SESSION_COMPLETED"; + const eventType = classifyEventType(topicToName(decodedTopics[0])); + + let rawPayloadObj: unknown; + if (decodedValue && typeof decodedValue === "object") { + rawPayloadObj = toPlainObject(decodedValue); + } else if (decodedValue !== undefined && decodedValue !== null) { + rawPayloadObj = { amount: toPlainObject(decodedValue) }; + } else { + rawPayloadObj = {}; } - const rawPayloadObj = - decodedValue && typeof decodedValue === "object" - ? (toPlainObject(decodedValue) as Record) + const payload: Record = Array.isArray(rawPayloadObj) + ? { data: rawPayloadObj } + : rawPayloadObj && typeof rawPayloadObj === "object" + ? { ...(rawPayloadObj as Record) } : {}; - const payload: Record = { ...rawPayloadObj }; + assignAliasedFields(payload); - // Standardise common Soroban snake_case keys to camelCase - if (payload["session_id"] !== undefined && !payload["sessionId"]) { - payload["sessionId"] = payload["session_id"]; - } - if (payload["seeker_address"] !== undefined && !payload["seekerAddress"]) { - payload["seekerAddress"] = payload["seeker_address"]; - } - if (payload["expert_address"] !== undefined && !payload["expertAddress"]) { - payload["expertAddress"] = payload["expert_address"]; - } - if (payload["expert_id"] !== undefined && !payload["expertId"]) { - payload["expertId"] = payload["expert_id"]; - } - if (payload["escrow_amount"] !== undefined && !payload["amount"]) { - payload["amount"] = payload["escrow_amount"]; + if (Array.isArray(payload["data"]) && payload["amount"] === undefined) { + const first = asScalar(payload["data"][0]); + if (first !== null) { + payload["amount"] = first; + } } - // Check if topic array carries positional args e.g. [topicName, sessionId] - if (!payload["sessionId"] && decodedTopics[1]) { - payload["sessionId"] = String(decodedTopics[1]); + if (!payload["sessionId"] && decodedTopics[1] != null) { + const sessionId = asScalar(toPlainObject(decodedTopics[1])); + if (sessionId !== null) { + payload["sessionId"] = String(sessionId); + } + } + if (!payload["amount"] && decodedTopics[2] != null) { + const amount = asScalar(toPlainObject(decodedTopics[2])); + if (amount !== null) { + payload["amount"] = amount; + } } return { eventType, payload }; @@ -144,31 +210,52 @@ export function decodeEventPayload( } } -/** - * Get the last processed ledger sequence from database state. - */ -export async function getLastProcessedLedger( +export async function getIndexerState( prisma: PrismaClient, key: string = INDEXER_KEY -): Promise { +): Promise<{ lastLedger: number; lastCursor: string | null }> { const state = await prisma.indexerState.findUnique({ where: { key }, }); - return state?.lastLedger ?? 0; + return { + lastLedger: state?.lastLedger ?? 0, + lastCursor: state?.lastCursor ?? null, + }; +} + +export async function getLastProcessedLedger( + prisma: PrismaClient, + key: string = INDEXER_KEY +): Promise { + const state = await getIndexerState(prisma, key); + return state.lastLedger; +} + +export async function getIndexerCursor( + prisma: PrismaClient, + key: string = INDEXER_KEY +): Promise { + const state = await getIndexerState(prisma, key); + return state.lastCursor; } -/** - * Update the last processed ledger sequence in database state. - */ export async function saveLastProcessedLedger( prisma: PrismaClient, lastLedger: number, - key: string = INDEXER_KEY + key: string = INDEXER_KEY, + lastCursor?: string | null ): Promise { await prisma.indexerState.upsert({ where: { key }, - create: { key, lastLedger }, - update: { lastLedger }, + create: { + key, + lastLedger, + lastCursor: lastCursor ?? null, + }, + update: { + lastLedger, + ...(lastCursor !== undefined ? { lastCursor } : {}), + }, }); } @@ -178,6 +265,33 @@ export interface SorobanIndexerOptions { server?: rpc.Server; } +export function buildContractEventFilters( + contractIds: string[] +): rpc.Api.EventFilter[] { + if (contractIds.length === 0) { + return [{ type: "contract" }]; + } + + const filters: rpc.Api.EventFilter[] = []; + for ( + let i = 0; + i < contractIds.length && filters.length < MAX_FILTERS; + i += MAX_CONTRACT_IDS_PER_FILTER + ) { + filters.push({ + type: "contract", + contractIds: contractIds.slice(i, i + MAX_CONTRACT_IDS_PER_FILTER), + }); + } + return filters; +} + +function ingestTxHash(event: rpc.Api.EventResponse): string { + if (event.txHash) return event.txHash; + if (event.id) return event.id; + throw new Error("Soroban event missing txHash and id"); +} + export class SorobanIndexerService { private prisma: PrismaClient; private server: rpc.Server; @@ -200,18 +314,20 @@ export class SorobanIndexerService { ] .filter(Boolean) .flatMap((id) => (id ? id.split(",") : [])) - .map((id) => id.trim()); + .map((id) => id.trim()) + .filter(Boolean); + + this.contractIds = options.contractIds ?? envContractIds; - this.contractIds = options.contractIds ?? (envContractIds.length > 0 ? envContractIds : []); + if (this.contractIds.length === 0) { + console.warn( + "[sorobanIndexer] No contract IDs configured; getEvents will match all contract events" + ); + } } /** - * Perform one polling tick: - * 1. Retrieve last processed ledger sequence. - * 2. Query Soroban RPC endpoint for contract events. - * 3. Decode event payloads safely. - * 4. Ingest into EventLog and execute processEvents (updating database Session/Transaction). - * 5. Update last processed ledger sequence. + * One polling tick: getEvents by contract ID → decode → EventLog → Session/Transaction. */ async pollOnce(): Promise<{ eventsFetched: number; @@ -224,82 +340,117 @@ export class SorobanIndexerService { this.isRunning = true; try { - let lastLedger = await getLastProcessedLedger(this.prisma); + const state = await getIndexerState(this.prisma); + let lastLedger = state.lastLedger; + const cursor = state.lastCursor; + let knownLatest: number | undefined; - if (lastLedger === 0) { + if (lastLedger === 0 && !cursor) { try { const latestRes = await this.server.getLatestLedger(); - lastLedger = Math.max(1, (latestRes.sequence ?? 1) - 10); + knownLatest = latestRes.sequence; + lastLedger = Math.max(1, (latestRes.sequence ?? 1) - LOOKBACK_LEDGERS); } catch (err) { - console.warn("[sorobanIndexer] Unable to fetch latest ledger, defaulting to ledger 1:", err); + console.warn( + "[sorobanIndexer] Unable to fetch latest ledger, defaulting to ledger 1:", + err + ); lastLedger = 1; } - await saveLastProcessedLedger(this.prisma, lastLedger); + await saveLastProcessedLedger(this.prisma, lastLedger, INDEXER_KEY, null); } - const startLedger = lastLedger + 1; + const filters = buildContractEventFilters(this.contractIds); + const collected: rpc.Api.EventResponse[] = []; + let newLastLedger = lastLedger; + let newCursor = cursor; + let startLedger = lastLedger + 1; - // Query Soroban RPC for contract events - let eventResponse: rpc.Api.GetEventsResponse; - try { - const filters: rpc.Api.EventFilter[] = [ - { - type: "contract", - ...(this.contractIds.length > 0 ? { contractIds: this.contractIds } : {}), - }, - ]; - - eventResponse = await this.server.getEvents({ - startLedger, - filters, - limit: 100, - }); - } catch (err) { - console.error("[sorobanIndexer] RPC getEvents error:", err); - return { eventsFetched: 0, processedCount: 0, latestLedger: lastLedger }; + if (!cursor) { + try { + const sequence = + knownLatest ?? (await this.server.getLatestLedger()).sequence; + if (typeof sequence === "number") { + startLedger = Math.min(Math.max(1, startLedger), sequence); + } + } catch { + // Keep lastLedger+1 when latest ledger is unavailable. + } } - const rawEvents = eventResponse.events ?? []; - let ingestedCount = 0; - let newLastLedger = lastLedger; + for (let page = 0; page < MAX_PAGES_PER_TICK; page++) { + let eventResponse: rpc.Api.GetEventsResponse; + try { + const request: rpc.Api.GetEventsRequest = newCursor + ? { filters, cursor: newCursor, limit: GET_EVENTS_LIMIT } + : { filters, startLedger, limit: GET_EVENTS_LIMIT }; - for (const rawEv of rawEvents) { - if (rawEv.ledger) { - newLastLedger = Math.max(newLastLedger, rawEv.ledger); + eventResponse = await this.server.getEvents(request); + } catch (err) { + console.error("[sorobanIndexer] RPC getEvents error:", err); + break; } - const topicRaw = (rawEv as unknown as { topic?: unknown[] }).topic ?? []; - const valueRaw = (rawEv as unknown as { value?: unknown }).value; - const txHash = (rawEv as unknown as { txHash?: string }).txHash ?? `tx_${rawEv.id}`; - - const { eventType, payload } = decodeEventPayload(topicRaw, valueRaw); - - if (eventType) { - const event: StellarEvent = { - txHash, - eventType, - payload: { - ...payload, - ledger: rawEv.ledger, - ledgerClosedAt: rawEv.ledgerClosedAt, - }, - }; - - const ingestRes = await ingestEvent(this.prisma, event); - if (ingestRes.created) { - ingestedCount++; + const rawEvents = eventResponse.events ?? []; + collected.push(...rawEvents); + + for (const rawEv of rawEvents) { + if (rawEv.ledger) { + newLastLedger = Math.max(newLastLedger, rawEv.ledger); } } + + if (eventResponse.cursor) { + newCursor = eventResponse.cursor; + } + + // Only fast-forward to RPC latest when this page is empty. A short + // non-empty page may still be a truncated scan; cursor continues it. + if (rawEvents.length === 0 && typeof eventResponse.latestLedger === "number") { + newLastLedger = Math.max(newLastLedger, eventResponse.latestLedger); + } + + if (rawEvents.length < GET_EVENTS_LIMIT) { + break; + } + } + + for (const rawEv of collected) { + const { eventType, payload } = decodeEventPayload( + rawEv.topic ?? [], + rawEv.value + ); + + if (!eventType) { + continue; + } + + const event: StellarEvent = { + txHash: ingestTxHash(rawEv), + eventType, + payload: { + ...payload, + ledger: rawEv.ledger, + ledgerClosedAt: rawEv.ledgerClosedAt, + }, + }; + + await ingestEvent(this.prisma, event); } const processRes = await processEvents(this.prisma); - if (newLastLedger > lastLedger) { - await saveLastProcessedLedger(this.prisma, newLastLedger); + if (newLastLedger > lastLedger || (newCursor && newCursor !== cursor)) { + await saveLastProcessedLedger( + this.prisma, + newLastLedger, + INDEXER_KEY, + newCursor + ); } return { - eventsFetched: rawEvents.length, + eventsFetched: collected.length, processedCount: processRes.processed, latestLedger: newLastLedger, }; diff --git a/backend/src/tests/eventListener.test.ts b/backend/src/tests/eventListener.test.ts index 036dd12..e9ac67a 100644 --- a/backend/src/tests/eventListener.test.ts +++ b/backend/src/tests/eventListener.test.ts @@ -417,6 +417,69 @@ describe("processEvents", () => { expect(result.errors[0]).toMatch(/missing amount/i); }); + it("processes PAYMENT_STREAMED into a transaction and keeps session ACTIVE", async () => { + if (!dbAvailable) return; + + const seeker = "GSEEKER_STREAM_EVT"; + const expertWallet = "GEXPERT_STREAM_EVT"; + await db.prisma.user.create({ data: { walletAddress: seeker } }); + const expertUser = await db.prisma.user.create({ + data: { walletAddress: expertWallet }, + }); + const expert = await db.prisma.expert.create({ + data: { userId: expertUser.id, name: "Stream Event Expert" }, + }); + await db.prisma.session.create({ + data: { + sessionId: "sess_stream_evt", + seekerAddress: seeker, + expertAddress: expertWallet, + expertId: expert.id, + status: "ACTIVE", + escrowAmount: 500n, + }, + }); + + await ingestEvent(db.prisma, { + txHash: "tx_stream_evt_001", + eventType: "PAYMENT_STREAMED", + payload: { + sessionId: "sess_stream_evt", + amount: "75", + ledgerClosedAt: "2026-08-22T12:00:00.000Z", + }, + }); + + const result = await processEvents(db.prisma); + expect(result.processed).toBe(1); + expect(result.errors).toHaveLength(0); + + const session = await db.prisma.session.findUnique({ + where: { sessionId: "sess_stream_evt" }, + }); + expect(session?.status).toBe("ACTIVE"); + + const tx = await db.prisma.transaction.findUnique({ + where: { txHash: "tx_stream_evt_001" }, + }); + expect(tx?.type).toBe("PAYMENT_RELEASED"); + expect(tx?.amount.toString()).toBe("75"); + expect(tx?.ledgerTime.toISOString()).toBe("2026-08-22T12:00:00.000Z"); + }); + + it("skips PAYMENT_STREAMED event with missing sessionId", async () => { + if (!dbAvailable) return; + await ingestEvent(db.prisma, { + txHash: "tx_stream_no_session", + eventType: "PAYMENT_STREAMED", + payload: { amount: 10 }, + }); + + const result = await processEvents(db.prisma); + expect(result.skipped).toBe(1); + expect(result.errors[0]).toMatch(/missing sessionId/i); + }); + it("notifies Discord/Telegram when a new SESSION_BOOKED session is created", async () => { if (!dbAvailable) return; diff --git a/backend/src/tests/sorobanIndexer.test.ts b/backend/src/tests/sorobanIndexer.test.ts index b3c06e1..695043c 100644 --- a/backend/src/tests/sorobanIndexer.test.ts +++ b/backend/src/tests/sorobanIndexer.test.ts @@ -3,6 +3,8 @@ import { parseScVal, toPlainObject, decodeEventPayload, + classifyEventType, + buildContractEventFilters, getLastProcessedLedger, saveLastProcessedLedger, SorobanIndexerService, @@ -75,6 +77,48 @@ describe("sorobanIndexer — XDR & Payload Decoding (Pure Unit Tests)", () => { expect(result).toBeDefined(); expect(result.eventType).toBeNull(); }); + + it("decodes PaymentStreamed topics and data without classifying as completion", () => { + const value = { + session_id: "sess-stream-1", + amount: "12345", + expert_address: "GEXPERT_STREAM", + }; + const { eventType, payload } = decodeEventPayload(["PaymentStreamed"], value); + + expect(eventType).toBe("PAYMENT_STREAMED"); + expect(classifyEventType("payment")).toBeNull(); + expect(payload["sessionId"]).toBe("sess-stream-1"); + expect(payload["amount"]).toBe("12345"); + expect(payload["expertAddress"]).toBe("GEXPERT_STREAM"); + }); + + it("reads sessionId from the second topic when the value map omits it", () => { + const { eventType, payload } = decodeEventPayload( + ["PaymentStreamed", "sess-from-topic"], + { amount: 50 } + ); + + expect(eventType).toBe("PAYMENT_STREAMED"); + expect(payload["sessionId"]).toBe("sess-from-topic"); + expect(payload["amount"]).toBe(50); + }); + + it("classifies only exact topic names, not substrings", () => { + expect(classifyEventType("PaymentStreamed")).toBe("PAYMENT_STREAMED"); + expect(classifyEventType("fund_session")).toBe("SESSION_BOOKED"); + expect(classifyEventType("booked")).toBe("SESSION_BOOKED"); + expect(classifyEventType("payment")).toBeNull(); + expect(classifyEventType("unbooked")).toBeNull(); + expect(classifyEventType("complete_session_v2")).toBeNull(); + }); + + it("builds getEvents filters scoped to contract IDs (max 5 per filter)", () => { + const filters = buildContractEventFilters(["C1", "C2"]); + expect(filters).toEqual([ + { type: "contract", contractIds: ["C1", "C2"] }, + ]); + }); }); describe("sorobanIndexer — Database & Service Logic", () => { @@ -182,6 +226,18 @@ describe("sorobanIndexer — Database & Service Logic", () => { expect(tx).not.toBeNull(); expect(tx?.type).toBe("ESCROW_FUNDED"); + expect(mockServer.getEvents).toHaveBeenCalledWith( + expect.objectContaining({ + startLedger: 100, + filters: [ + expect.objectContaining({ + type: "contract", + contractIds: ["CCONTRACT123"], + }), + ], + }) + ); + const lastLedger = await getLastProcessedLedger(db.prisma); expect(lastLedger).toBe(100); }); @@ -260,4 +316,112 @@ describe("sorobanIndexer — Database & Service Logic", () => { }); expect(tx?.type).toBe("PAYMENT_RELEASED"); }); + + it("indexes PaymentStreamed into Transaction without completing the session", async () => { + if (!dbAvailable) return; + + await db.prisma.user.create({ data: { walletAddress: "GSEEKER_STREAM" } }); + const expUser = await db.prisma.user.create({ + data: { walletAddress: "GEXPERT_STREAM" }, + }); + const exp = await db.prisma.expert.create({ + data: { id: "exp_stream", userId: expUser.id, name: "Stream Expert" }, + }); + await db.prisma.session.create({ + data: { + sessionId: "sess-stream-live", + seekerAddress: "GSEEKER_STREAM", + expertAddress: "GEXPERT_STREAM", + expertId: exp.id, + status: "ACTIVE", + escrowAmount: 1_000_000n, + }, + }); + + const ledgerClosedAt = "2026-08-22T18:00:00.000Z"; + const mockEvents = [ + { + id: "0000003000-0000000001", + ledger: 300, + ledgerClosedAt, + contractId: "CESCROW", + topic: ["PaymentStreamed", "sess-stream-live"], + value: { amount: "250000" }, + txHash: "tx_payment_streamed_001", + }, + ]; + + const mockServer = { + getLatestLedger: jest.fn().mockResolvedValue({ sequence: 305 }), + getEvents: jest.fn().mockResolvedValue({ + events: mockEvents, + cursor: "cursor-300", + latestLedger: 305, + }), + } as unknown as any; + + const service = new SorobanIndexerService(db.prisma, { + server: mockServer, + contractIds: ["CESCROW"], + }); + + await saveLastProcessedLedger(db.prisma, 299); + + const started = Date.now(); + const pollRes = await service.pollOnce(); + const elapsed = Date.now() - started; + + expect(elapsed).toBeLessThan(2000); + expect(pollRes.eventsFetched).toBe(1); + expect(pollRes.processedCount).toBe(1); + + const session = await db.prisma.session.findUnique({ + where: { sessionId: "sess-stream-live" }, + }); + expect(session?.status).toBe("ACTIVE"); + + const tx = await db.prisma.transaction.findUnique({ + where: { txHash: "tx_payment_streamed_001" }, + }); + expect(tx).not.toBeNull(); + expect(tx?.type).toBe("PAYMENT_RELEASED"); + expect(tx?.amount.toString()).toBe("250000"); + expect(tx?.sessionId).toBe("sess-stream-live"); + expect(tx?.ledgerTime.toISOString()).toBe(ledgerClosedAt); + + expect(mockServer.getEvents).toHaveBeenCalledWith( + expect.objectContaining({ + filters: [ + expect.objectContaining({ + type: "contract", + contractIds: ["CESCROW"], + }), + ], + }) + ); + + const state = await db.prisma.indexerState.findUnique({ + where: { key: "stellar_soroban_indexer" }, + }); + expect(state?.lastCursor).toBe("cursor-300"); + }); + + it("resumes getEvents from the saved cursor instead of startLedger", async () => { + if (!dbAvailable) return; + await saveLastProcessedLedger(db.prisma, 400, "stellar_soroban_indexer", "cursor-400"); + + const getEvents = jest.fn().mockResolvedValue({ events: [], cursor: "cursor-400" }); + const mockServer = { + getLatestLedger: jest.fn().mockResolvedValue({ sequence: 410 }), + getEvents, + } as unknown as any; + + const service = new SorobanIndexerService(db.prisma, { server: mockServer }); + await service.pollOnce(); + + expect(getEvents).toHaveBeenCalledWith( + expect.objectContaining({ cursor: "cursor-400", limit: 100 }) + ); + expect(getEvents.mock.calls[0][0].startLedger).toBeUndefined(); + }); }); diff --git a/docker-compose.yml b/docker-compose.yml index db99503..6f554ab 100644 --- a/docker-compose.yml +++ b/docker-compose.yml @@ -54,7 +54,9 @@ services: environment: NODE_ENV: production DATABASE_URL: postgresql://skillsphere:skillsphere@db:5432/skillsphere?schema=public - INDEXER_INTERVAL_MS: "5000" + INDEXER_INTERVAL_MS: "1000" + SOROBAN_RPC_URL: "https://soroban-testnet.stellar.org" + ESCROW_CONTRACT_ID: ${ESCROW_CONTRACT_ID:-} depends_on: db: condition: service_healthy