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
68 changes: 68 additions & 0 deletions backend/SETUP_GUIDE.md
Original file line number Diff line number Diff line change
Expand Up @@ -101,11 +101,79 @@ MONITORING_PROCESSING_TIME_THRESHOLD=30000
MONITORING_GAS_FEE_SPIKE_THRESHOLD=2.0
MONITORING_NETWORK_CONGESTION_THRESHOLD=0.8

# ─────────────────────────────────────────────────────────────────
# Event Indexer Configuration (NEW)
# ─────────────────────────────────────────────────────────────────
# Master switch – set to "true" to enable the on-chain event indexer.
# When disabled the service starts normally but no Soroban events are polled.
EVENT_INDEXER_ENABLED=false

# Soroban RPC endpoint the indexer uses to fetch events.
# Defaults to the public Stellar testnet RPC if not set.
SOROBAN_RPC_URL=https://soroban-testnet.stellar.org

# Comma-separated list of Soroban contract addresses to watch.
# Example: CAABC...XYZ,CBBDE...QRS
# Leave empty to disable indexing even when EVENT_INDEXER_ENABLED=true.
EVENT_INDEXER_CONTRACT_IDS=

# How often (in milliseconds) the indexer polls for new events.
# Default: 5000 (5 seconds). Minimum recommended: 2000.
EVENT_INDEXER_POLL_INTERVAL=5000

# Number of ledgers to process per poll batch.
# Larger values reduce HTTP round-trips but increase per-batch latency.
# Default: 100
EVENT_INDEXER_BATCH_SIZE=100

# Optional: override the starting ledger for back-fills or initial sync.
# If not set the indexer resumes from the last saved checkpoint (or ledger 0).
# EVENT_INDEXER_START_LEDGER=

# Logging
LOG_LEVEL=info
LOG_FILE_PATH=./logs
```

#### Event Indexer Quick Start

1. Set `EVENT_INDEXER_ENABLED=true` in your `.env`
2. Set `EVENT_INDEXER_CONTRACT_IDS` to your deployed Soroban contract address(es)
3. Confirm `SOROBAN_RPC_URL` points to the right network (testnet or mainnet)
4. Run the database migration to create the required tables:
```bash
npm run migrate:up
```
5. Start the server – the indexer will boot automatically:
```bash
npm run dev
```
6. Verify via the health endpoint:
```bash
curl http://localhost:3001/health | jq '.eventIndexer'
# Expected: {"status":"running","lastLedger":12345,"eventsProcessed":42,"lag":3}
```

#### Admin API (requires admin JWT)

| Method | Path | Description |
|--------|------|-------------|
| `GET` | `/api/v1/indexer/status` | Current indexer status |
| `POST` | `/api/v1/indexer/start` | Start the indexer |
| `POST` | `/api/v1/indexer/stop` | Gracefully stop the indexer |

#### Indexed Event Types

| Event Type | Soroban topic key | Domain side-effect |
|---|---|---|
| `CredentialIssued` | `cred:issued` | Upserts row in `credentials` table |
| `CredentialRevoked` | `cred:revoked` | Sets `status='revoked'` in `credentials` |
| `CourseCreated` | `course:created` | Inserts row in `courses` table |
| `EnrollmentCreated` | `enroll:created` | Inserts row in `enrollments` table |
| `AchievementMinted` | `ach:minted` | Logged only |
| `PaymentReceived` | `pay:received` | Logged only |
| `ProfileUpdated` | `profile:update` | Logged only |

### 4. Redis Setup

#### Option 1: Local Redis Installation
Expand Down
55 changes: 55 additions & 0 deletions backend/migrations/003_create_indexed_events.sql
Original file line number Diff line number Diff line change
@@ -0,0 +1,55 @@
-- UP
-- Migration: Create indexed_events table for Stellar/Soroban contract event indexing
-- Each event is uniquely keyed by (contract_id, ledger, event_index) to prevent duplicates.

CREATE TABLE IF NOT EXISTS indexed_events (
id UUID PRIMARY KEY DEFAULT gen_random_uuid(),
contract_id VARCHAR(64) NOT NULL,
ledger BIGINT NOT NULL,
event_index INTEGER NOT NULL,
event_type VARCHAR(64) NOT NULL,
topic TEXT[] NOT NULL DEFAULT '{}',
payload JSONB NOT NULL DEFAULT '{}',
processed_at TIMESTAMP WITH TIME ZONE NOT NULL DEFAULT NOW(),
created_at TIMESTAMP WITH TIME ZONE NOT NULL DEFAULT NOW(),

-- Deduplication: same event on same contract at same ledger position must not appear twice
CONSTRAINT uq_indexed_events_key UNIQUE (contract_id, ledger, event_index)
);

-- Index for efficient range queries by ledger (used by the indexer resume logic)
CREATE INDEX IF NOT EXISTS idx_indexed_events_ledger
ON indexed_events (ledger ASC);

-- Index for querying by contract
CREATE INDEX IF NOT EXISTS idx_indexed_events_contract_id
ON indexed_events (contract_id);

-- Index for querying by event type
CREATE INDEX IF NOT EXISTS idx_indexed_events_event_type
ON indexed_events (event_type);

-- Index to support last-processed-ledger checkpoint lookups
CREATE INDEX IF NOT EXISTS idx_indexed_events_processed_at
ON indexed_events (processed_at);

-- Table to persist the indexer checkpoint (last successfully indexed ledger per contract set)
CREATE TABLE IF NOT EXISTS indexer_checkpoints (
id SERIAL PRIMARY KEY,
checkpoint_key VARCHAR(64) NOT NULL UNIQUE, -- e.g. 'default' or contract group name
last_ledger BIGINT NOT NULL DEFAULT 0,
updated_at TIMESTAMP WITH TIME ZONE NOT NULL DEFAULT NOW()
);

-- Seed a default checkpoint row so the indexer can always UPDATE rather than INSERT/UPDATE
INSERT INTO indexer_checkpoints (checkpoint_key, last_ledger)
VALUES ('default', 0)
ON CONFLICT (checkpoint_key) DO NOTHING;

-- @undo
DROP TABLE IF EXISTS indexer_checkpoints;
DROP INDEX IF EXISTS idx_indexed_events_processed_at;
DROP INDEX IF EXISTS idx_indexed_events_event_type;
DROP INDEX IF EXISTS idx_indexed_events_contract_id;
DROP INDEX IF EXISTS idx_indexed_events_ledger;
DROP TABLE IF EXISTS indexed_events;
85 changes: 71 additions & 14 deletions backend/src/index.js
Original file line number Diff line number Diff line change
Expand Up @@ -19,6 +19,10 @@ const transactionProcessor = require('./workers/transactionProcessor');
const transactionEvents = require('./events/transactionEvents');
const emailWorker = require('./workers/emailWorker');

// Event Indexer – polls Soroban for on-chain events and syncs them to PostgreSQL
let eventIndexerInstance = null;
const EVENT_INDEXER_ENABLED = process.env.EVENT_INDEXER_ENABLED === 'true';

// Import security middleware
const {
securityPerformanceTracker,
Expand Down Expand Up @@ -216,22 +220,46 @@ v1Router.use('/cross-protocol-bridge', crossProtocolBridgeRoutes);
const adminRoutes = require('./routes/admin');
v1Router.use('/admin', adminRoutes);

// Schemas helper for versioned responses
const { createVersionedResponse } = require('./utils/schemas');
const { errorHandler } = require('./middleware/errorHandler');
const { ValidationError } = require('./utils/errors');
const { getCompressionStats } = require('./middleware/compression');
// Event Indexer admin routes (start / stop / status)
const indexerAdminRouter = require('express').Router();

// Health check under v1 for OpenAPI spec compatibility
v1Router.get('/health', (req, res) => {
const version = req.apiVersion || 'v1';
res.json(createVersionedResponse({
status: 'healthy',
uptime: process.uptime(),
supportedVersions: SUPPORTED_VERSIONS,
}, version));
indexerAdminRouter.get('/status', (req, res) => {
try {
const { getIndexerStatus } = require('./services/eventIndexer');
res.json({ eventIndexer: getIndexerStatus() });
} catch (err) {
res.json({ eventIndexer: { status: 'stopped', error: err.message } });
}
});

indexerAdminRouter.post('/start', async (req, res) => {
try {
if (!eventIndexerInstance) {
return res.status(400).json({ error: 'Indexer not initialized' });
}
await eventIndexerInstance.start();
const { getIndexerStatus } = require('./services/eventIndexer');
res.json({ message: 'Indexer started', status: getIndexerStatus() });
} catch (err) {
res.status(500).json({ error: err.message });
}
});

indexerAdminRouter.post('/stop', async (req, res) => {
try {
if (!eventIndexerInstance) {
return res.status(400).json({ error: 'Indexer not initialized' });
}
await eventIndexerInstance.stop();
const { getIndexerStatus } = require('./services/eventIndexer');
res.json({ message: 'Indexer stopped', status: getIndexerStatus() });
} catch (err) {
res.status(500).json({ error: err.message });
}
});

v1Router.use('/indexer', require('./middleware/auth').requireAdmin, indexerAdminRouter);

// Mount v1 router at /api/v1
app.use('/api/v1', v1Router);

Expand Down Expand Up @@ -304,6 +332,27 @@ async function startServer() {
await transactionEvents.startListening();
emailWorker.getEmailWorker().start();

// Start the event indexer if enabled
if (EVENT_INDEXER_ENABLED) {
try {
const { Pool } = require('pg');
const { getEventIndexer } = require('./services/eventIndexer');
const indexerPool = new Pool({
connectionString: process.env.DATABASE_URL || 'postgresql://postgres:postgres@localhost:5432/starked',
max: 5, // dedicated small pool for the indexer
idleTimeoutMillis: 30000,
connectionTimeoutMillis: 5000,
});
eventIndexerInstance = getEventIndexer(indexerPool);
await eventIndexerInstance.start();
console.log('🔗 Event Indexer started – polling Soroban for on-chain events');
} catch (indexerErr) {
console.error('⚠️ Event Indexer failed to start (non-fatal):', indexerErr.message);
}
} else {
console.log('ℹ️ Event Indexer disabled. Set EVENT_INDEXER_ENABLED=true to enable.');
}

server.listen(PORT, () => {
console.log(`🚀 StarkEd Education Backend running on port ${PORT}`);
console.log(`📚 Quiz Management API available at /api/v1/quizzes`);
Expand All @@ -317,6 +366,7 @@ async function startServer() {
console.log(`🌐 Federated Learning API available at /api/v1/federated-learning`);
console.log(`🧠 AGI Tutor API available at /api/v1/agi-tutor`);
console.log(`🔐 Quantum-Resistant Secure Communication API available at /api/v1/secure-comm`);
console.log(`🔗 Event Indexer API available at /api/v1/indexer (admin-only)`);
console.log(`🏥 Health check available at /api/health`);
console.log(`✅ Transaction Queue System initialized successfully`);
});
Expand All @@ -328,7 +378,14 @@ async function startServer() {

process.on('SIGINT', async () => {
console.log('SIGINT received, shutting down gracefully...');
emailWorker.getEmailWorker().stop();
if (eventIndexerInstance) {
try {
await eventIndexerInstance.stop();
console.log('Event Indexer stopped cleanly.');
} catch (err) {
console.error('Error stopping event indexer:', err.message);
}
}
await transactionQueue.stopProcessing();
await transactionProcessor.stop();
await transactionEvents.stopListening();
Expand Down
Loading