Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
8 changes: 6 additions & 2 deletions backend/.env.example
Original file line number Diff line number Diff line change
Expand Up @@ -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
Expand Down
Original file line number Diff line number Diff line change
@@ -0,0 +1,2 @@
-- AlterTable
ALTER TABLE "IndexerState" ADD COLUMN "lastCursor" TEXT;
4 changes: 3 additions & 1 deletion backend/prisma/schema.prisma
Original file line number Diff line number Diff line change
Expand Up @@ -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
Expand All @@ -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
}

130 changes: 107 additions & 23 deletions backend/src/eventListener.ts
Original file line number Diff line number Diff line change
@@ -1,4 +1,4 @@
import { PrismaClient } from "@prisma/client";
import { Prisma, PrismaClient } from "@prisma/client";
import { getNotificationService } from "./notificationService";
import {
publishSessionStatus,
Expand All @@ -9,6 +9,7 @@ export type StellarEventType =
| "SESSION_BOOKED"
| "SESSION_COMPLETED"
| "PAYMENT_RELEASED"
| "PAYMENT_STREAMED"
| "EXPERT_REGISTERED"
| "SESSION_PAUSED"
| "SESSION_REFUNDED";
Expand All @@ -18,6 +19,7 @@ const SESSION_SCOPED_EVENTS: ReadonlySet<StellarEventType> = new Set([
"SESSION_BOOKED",
"SESSION_COMPLETED",
"PAYMENT_RELEASED",
"PAYMENT_STREAMED",
"SESSION_PAUSED",
"SESSION_REFUNDED",
]);
Expand Down Expand Up @@ -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;
}
}

/**
Expand Down Expand Up @@ -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;
Expand Down Expand Up @@ -249,17 +268,16 @@ 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: {
txHash: hash,
sessionId,
amount: escrowAmount,
type: "ESCROW_FUNDED",
ledgerTime,
ledgerTime: ledgerTimeFromPayload(payload),
},
update: {},
});
Expand All @@ -279,19 +297,18 @@ 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: {
txHash: hash,
sessionId,
amount,
type: "PAYMENT_RELEASED",
ledgerTime,
ledgerTime: ledgerTimeFromPayload(payload),
},
update: {},
});
Expand All @@ -317,23 +334,67 @@ 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: {
txHash: hash,
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<string, unknown>,
txHash?: string
): Promise<void> {
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<string, unknown>
Expand All @@ -360,25 +421,46 @@ 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: {
txHash: hash,
sessionId,
amount,
type: "REFUND_ISSUED",
ledgerTime,
ledgerTime: ledgerTimeFromPayload(payload),
},
update: {},
});
}
}

function transactionHash(
payload: Record<string, unknown>,
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<string, unknown>): 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.
Expand All @@ -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(
Expand Down
27 changes: 17 additions & 10 deletions backend/src/indexer.ts
Original file line number Diff line number Diff line change
Expand Up @@ -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<void> {
// 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) {
Expand All @@ -32,11 +42,8 @@ async function tick(): Promise<void> {
}

async function main(): Promise<void> {
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");

Expand Down
Loading