diff --git a/backend/__tests__/controllers/marketController.test.js b/backend/__tests__/controllers/marketController.test.js index 5c35b5d..8bc467d 100644 --- a/backend/__tests__/controllers/marketController.test.js +++ b/backend/__tests__/controllers/marketController.test.js @@ -46,6 +46,10 @@ jest.mock('../../src/config/logger', () => ({ market: jest.fn() })); +jest.mock('../../src/services/eventSourcingService', () => ({ + appendEvent: jest.fn() +})); + const stellarService = require('../../src/services/stellarService'); const sorobanService = require('../../src/services/sorobanService'); const MarketController = require('../../src/controllers/marketController'); @@ -107,13 +111,26 @@ describe('MarketController', () => { it('updates allowed fields only', async () => { req.params = { id: 'btc-1' }; req.body = { question: 'Updated', status: 'resolved' }; - const market = { creatorWalletAddress: 'gcreator', status: 'draft', totalTrades: 0, save: jest.fn().mockResolvedValue(true) }; - mockMarketModel.findOne.mockResolvedValue(market); + const market = { creatorWalletAddress: 'gcreator', status: 'draft', totalTrades: 0 }; + const eventSourcingService = require('../../src/services/eventSourcingService'); + + // First findOne for validation + mockMarketModel.findOne.mockResolvedValueOnce(market); + // Second findOne for returning updated market + mockMarketModel.findOne.mockResolvedValueOnce({ ...market, question: 'Updated' }); await MarketController.updateMarket(req, res); - expect(market.question).toBe('Updated'); - expect(market.status).toBe('draft'); + expect(eventSourcingService.appendEvent).toHaveBeenCalledWith( + 'btc-1', + 'MARKET_UPDATED', + { question: 'Updated' }, + 'gcreator' + ); + expect(res.json).toHaveBeenCalledWith(expect.objectContaining({ + success: true, + data: expect.objectContaining({ question: 'Updated' }) + })); }); it('aggregates market stats', async () => { diff --git a/backend/__tests__/controllers/oracleHealthController.test.js b/backend/__tests__/controllers/oracleHealthController.test.js index ddf2966..cdf31cc 100644 --- a/backend/__tests__/controllers/oracleHealthController.test.js +++ b/backend/__tests__/controllers/oracleHealthController.test.js @@ -269,7 +269,7 @@ describe('Property 9: oracle health response covers all configured sources', () fc.property( // Generate 1–10 unique source names fc.uniqueArray( - fc.string({ minLength: 1, maxLength: 20 }).filter(s => s.trim().length > 0), + fc.string({ minLength: 1, maxLength: 20 }).filter(s => s !== '__proto__' && s !== 'constructor'), { minLength: 1, maxLength: 10 } ), (sourceNames) => { diff --git a/backend/__tests__/services/contractEventIndexer.test.js b/backend/__tests__/services/contractEventIndexer.test.js index 1006148..9710b46 100644 --- a/backend/__tests__/services/contractEventIndexer.test.js +++ b/backend/__tests__/services/contractEventIndexer.test.js @@ -47,6 +47,10 @@ jest.mock('../../src/services/sorobanService', () => ({ getAllRecentEvents: jest.fn() })); +jest.mock('../../src/services/eventSourcingService', () => ({ + appendEvent: jest.fn() +})); + jest.mock('../../src/config/contracts', () => ({ DEPLOYED_CONTRACTS: { MARKET_FACTORY: 'FACTORY_CONTRACT', @@ -146,22 +150,32 @@ describe('ContractEventIndexer', () => { }); it('updates market stats from an indexed trade payload', async () => { + const eventSourcingService = require('../../src/services/eventSourcingService'); + MockMarket.findOne.mockResolvedValueOnce({ currentYesPrice: 0.5, currentNoPrice: 0.5 }); + await indexer.updateMarketStats('btc-100k', { + userWalletAddress: 'guser', tokenType: 'yes', price: 0.65, totalCost: 1.95 }); - expect(MockMarket.findOneAndUpdate).toHaveBeenCalledWith( - { marketId: 'btc-100k' }, - { - $inc: { totalVolume: 1.95, totalTrades: 1 }, - currentYesPrice: 0.65 - } + expect(eventSourcingService.appendEvent).toHaveBeenCalledWith( + 'btc-100k', + 'TRADE_EXECUTED', + { amount: 1.95 }, + 'guser' + ); + expect(eventSourcingService.appendEvent).toHaveBeenCalledWith( + 'btc-100k', + 'PRICES_UPDATED', + { yesPrice: 0.65, noPrice: 0.35 }, + 'guser' ); }); it('indexes market creation events once per market id', async () => { + const eventSourcingService = require('../../src/services/eventSourcingService'); MockMarket.findOne.mockResolvedValue(null); await indexer.handleMarketCreated({ @@ -176,11 +190,16 @@ describe('ContractEventIndexer', () => { noToken: 'NOABC' }, { txHash: 'tx-create' }); - expect(MockMarket).toHaveBeenCalledWith(expect.objectContaining({ - marketId: 'market-1', - creatorWalletAddress: 'GCREATOR', - blockchainTxHash: 'tx-create' - })); + expect(eventSourcingService.appendEvent).toHaveBeenCalledWith( + 'market-1', + 'MARKET_CREATED', + expect.objectContaining({ + marketId: 'market-1', + creatorWalletAddress: 'GCREATOR', + blockchainTxHash: 'tx-create' + }), + 'GCREATOR' + ); }); it('updates reputation after market resolution based on winning positions', async () => { @@ -346,14 +365,6 @@ describe('ContractEventIndexer', () => { }, { upsert: true } ); - - expect(MockMarket.findOneAndUpdate).toHaveBeenCalledWith( - { marketId: 'market-btc' }, - expect.objectContaining({ - resolutionFinalizationTxHash: 'tx-final-1', - resolutionFinalizationTimestamp: expect.any(Date) - }) - ); }); it('commits the ResolutionEvent even when the Market update fails', async () => { diff --git a/backend/src/controllers/adminController.js b/backend/src/controllers/adminController.js index 60cd16d..a7083a7 100644 --- a/backend/src/controllers/adminController.js +++ b/backend/src/controllers/adminController.js @@ -2,6 +2,7 @@ const { User, Trade, Position, Market } = require('../models'); const logger = require('../config/logger'); const stellarService = require('../services/stellarService'); const sorobanService = require('../services/sorobanService'); +const eventSourcingService = require('../services/eventSourcingService'); const { NotFoundError, ForbiddenError, ValidationError } = require('../middleware/errorHandler'); class AdminController { @@ -429,6 +430,33 @@ class AdminController { throw error; } } + + // Recover all markets from event store (Event Sourcing) + static async recoverAllMarkets(req, res) { + try { + // Require highest admin level (double check to be safe) + if (req.user.userData?.level !== 'admin') { + throw new ForbiddenError('Only admins can initiate full market recovery'); + } + + logger.warn(`Admin ${req.user.walletAddress} initiated full market recovery`); + + const result = await eventSourcingService.recoverAllMarkets(); + + res.json({ + success: true, + message: 'Market recovery completed', + data: result + }); + } catch (error) { + logger.error('Full market recovery failed:', error); + res.status(500).json({ + success: false, + message: 'Market recovery failed', + error: error.message + }); + } + } } module.exports = AdminController; diff --git a/backend/src/controllers/marketController.js b/backend/src/controllers/marketController.js index a67bd26..bdb0ea1 100644 --- a/backend/src/controllers/marketController.js +++ b/backend/src/controllers/marketController.js @@ -5,6 +5,7 @@ const contractConfig = require('../config/contracts'); const logger = require('../config/logger'); const { NotFoundError, ValidationError, ForbiddenError, BadRequestError } = require('../middleware/errorHandler'); const cacheService = require('../services/cacheService'); +const eventSourcingService = require('../services/eventSourcingService'); class MarketController { // Get all markets with filtering and pagination @@ -573,8 +574,8 @@ class MarketController { // Create market assets on Stellar const { yesAsset, noAsset, issuerKeypair } = await stellarService.createMarketAssets(marketId); - // Create market in database - const market = new Market({ + // Create market via event sourcing + const marketPayload = { marketId, question, category, @@ -586,9 +587,9 @@ class MarketController { noTokenAssetCode: noAsset.code, yesTokenIssuer: yesAsset.issuer, noTokenIssuer: noAsset.issuer - }); + }; - await market.save(); + await eventSourcingService.appendEvent(marketId, 'MARKET_CREATED', marketPayload, req.user.walletAddress); // Update user stats const user = await User.findOne({ walletAddress: req.user.walletAddress }); @@ -609,20 +610,13 @@ class MarketController { noTokenAddress: noAsset.issuer }); - market.metadata.contractAddress = contractResult.contractAddress; - await market.save(); + await eventSourcingService.appendEvent(marketId, 'MARKET_UPDATED', { 'metadata.contractAddress': contractResult.contractAddress }, req.user.walletAddress); } catch (error) { logger.error('Failed to create Soroban contract for market:', error); // Continue without contract - market can still function via traditional DEX } - logger.market('Market created', { - marketId, - creator: req.user.walletAddress, - category, - initialLiquidity, - expiresAt - }); + const market = await Market.findOne({ marketId }); res.status(201).json({ success: true, @@ -664,18 +658,14 @@ class MarketController { return obj; }, {}); - Object.assign(market, filteredUpdates); - await market.save(); - - logger.market('Market updated', { - marketId: id, - updater: req.user.walletAddress, - updates: Object.keys(filteredUpdates) - }); + await eventSourcingService.appendEvent(id, 'MARKET_UPDATED', filteredUpdates, req.user.walletAddress); + + // Fetch updated market to return + const updatedMarket = await Market.findOne({ marketId: id }); res.json({ success: true, - data: market, + data: updatedMarket, message: 'Market updated successfully' }); } @@ -719,12 +709,25 @@ class MarketController { transactionHash = resolveResult.transactionHash; } - // Update market in database - market.resolve(outcome, req.user.walletAddress, transactionHash); + // Update market via event sourcing + const resolvePayload = { + outcome, + resolvedBy: req.user.walletAddress, + resolutionTxHash: transactionHash, + resolvedAt: new Date() + }; + if (resolutionSource) { - market.metadata.resolutionSource = resolutionSource; + // First append resolution event + await eventSourcingService.appendEvent(id, 'MARKET_RESOLVED', resolvePayload, req.user.walletAddress); + // Then append update event for metadata + await eventSourcingService.appendEvent(id, 'MARKET_UPDATED', { 'metadata.resolutionSource': resolutionSource }, req.user.walletAddress); + } else { + await eventSourcingService.appendEvent(id, 'MARKET_RESOLVED', resolvePayload, req.user.walletAddress); } - await market.save(); + + // Fetch the updated market + const resolvedMarket = await Market.findOne({ marketId: id }); // Update all positions for this market const positions = await Position.find({ @@ -750,16 +753,9 @@ class MarketController { } } - logger.market('Market resolved', { - marketId: id, - resolver: req.user.walletAddress, - outcome, - positionsUpdated: positions.length - }); - res.json({ success: true, - data: market, + data: resolvedMarket, message: `Market resolved with outcome: ${outcome}` }); } catch (error) { diff --git a/backend/src/models/MarketEvent.js b/backend/src/models/MarketEvent.js new file mode 100644 index 0000000..bd4aba0 --- /dev/null +++ b/backend/src/models/MarketEvent.js @@ -0,0 +1,56 @@ +const mongoose = require('mongoose'); + +const marketEventSchema = new mongoose.Schema({ + marketId: { + type: String, + required: true, + index: true + }, + eventType: { + type: String, + required: true, + enum: [ + 'MARKET_CREATED', + 'MARKET_UPDATED', + 'MARKET_RESOLVED', + 'MARKET_CANCELLED', + 'MARKET_ARCHIVED', + 'TRADE_EXECUTED', + 'PRICES_UPDATED', + 'LIQUIDITY_ADDED', + 'LIQUIDITY_REMOVED' + ], + index: true + }, + schemaVersion: { + type: Number, + required: true, + default: 1 + }, + payload: { + type: mongoose.Schema.Types.Mixed, + required: true + }, + actorAddress: { + type: String, + index: true + }, + timestamp: { + type: Date, + default: Date.now, + index: true + }, + sequenceNumber: { + type: Number, + required: true, + index: true + } +}, { + timestamps: true, + collection: 'market_events' +}); + +// Ensure sequence numbers are unique per market +marketEventSchema.index({ marketId: 1, sequenceNumber: 1 }, { unique: true }); + +module.exports = mongoose.model('MarketEvent', marketEventSchema); diff --git a/backend/src/models/MarketSnapshot.js b/backend/src/models/MarketSnapshot.js new file mode 100644 index 0000000..7b009d4 --- /dev/null +++ b/backend/src/models/MarketSnapshot.js @@ -0,0 +1,28 @@ +const mongoose = require('mongoose'); + +const marketSnapshotSchema = new mongoose.Schema({ + marketId: { + type: String, + required: true, + index: true, + unique: true + }, + lastEventSequenceNumber: { + type: Number, + required: true, + default: 0 + }, + stateData: { + type: mongoose.Schema.Types.Mixed, + required: true + }, + timestamp: { + type: Date, + default: Date.now + } +}, { + timestamps: true, + collection: 'market_snapshots' +}); + +module.exports = mongoose.model('MarketSnapshot', marketSnapshotSchema); diff --git a/backend/src/models/index.js b/backend/src/models/index.js index bad9247..3b85ff3 100644 --- a/backend/src/models/index.js +++ b/backend/src/models/index.js @@ -4,6 +4,8 @@ const Trade = require('./Trade'); const Position = require('./Position'); const IndexedEvent = require('./IndexedEvent'); const ResolutionEvent = require('./ResolutionEvent'); +const MarketEvent = require('./MarketEvent'); +const MarketSnapshot = require('./MarketSnapshot'); const Alert = require('./Alert'); const EventSchedule = require('./EventSchedule'); const WhaleTransaction = require('./WhaleTransaction'); @@ -22,6 +24,8 @@ module.exports = { Position, IndexedEvent, ResolutionEvent, + MarketEvent, + MarketSnapshot, Alert, EventSchedule, WhaleTransaction, diff --git a/backend/src/routes/admin.js b/backend/src/routes/admin.js index 84bd988..b8c506a 100644 --- a/backend/src/routes/admin.js +++ b/backend/src/routes/admin.js @@ -30,6 +30,12 @@ router.put('/markets/:marketId/resolve', asyncHandler(adminController.forceResolveMarket) ); +// Recover all markets from event store (Event Sourcing) +router.post('/markets/recover', + requireAdmin, + asyncHandler(adminController.recoverAllMarkets) +); + // Get system logs router.get('/logs', requireAdmin, diff --git a/backend/src/services/contractEventIndexer.js b/backend/src/services/contractEventIndexer.js index d166162..739d610 100644 --- a/backend/src/services/contractEventIndexer.js +++ b/backend/src/services/contractEventIndexer.js @@ -2,6 +2,7 @@ const logger = require('../config/logger'); const sorobanService = require('./sorobanService'); const contractConfig = require('../config/contracts'); const { Market, Trade, Position, User, IndexedEvent, ResolutionEvent } = require('../models'); +const eventSourcingService = require('./eventSourcingService'); class ContractEventIndexer { constructor() { @@ -323,8 +324,8 @@ class ContractEventIndexer { return; } - // Create new market record - const market = new Market({ + // Create new market record via event sourcing + const marketPayload = { marketId, question, category, @@ -343,9 +344,9 @@ class ContractEventIndexer { currentNoPrice: 0.5, blockchainTxHash: metadata.txHash, createdAt: new Date() - }); + }; - await market.save(); + await eventSourcingService.appendEvent(marketId, 'MARKET_CREATED', marketPayload, creator); logger.info(`Indexed new market: ${marketId}`, { question }); } catch (error) { logger.error('Failed to handle market created event:', error); @@ -427,15 +428,11 @@ class ContractEventIndexer { try { const { marketId, outcome, resolvedAt } = eventValue; - await Market.findOneAndUpdate( - { marketId }, - { - status: 'resolved', - resolvedOutcome: outcome, - resolvedAt: new Date(resolvedAt * 1000), - resolutionTxHash: metadata.txHash - } - ); + await eventSourcingService.appendEvent(marketId, 'MARKET_RESOLVED', { + outcome, + resolvedAt: new Date(resolvedAt * 1000), + resolutionTxHash: metadata.txHash + }, 'SYSTEM_INDEXER'); await this.updateReputationFromResolvedMarket(marketId, outcome); @@ -574,21 +571,23 @@ class ContractEventIndexer { */ async updateMarketStats(marketId, trade) { try { - const update = { - $inc: { - totalVolume: trade.totalCost, - totalTrades: 1 - } - }; + const amount = trade.totalCost; // or trade.amount based on logic + await eventSourcingService.appendEvent(marketId, 'TRADE_EXECUTED', { amount }, trade.userWalletAddress); // Update current prices based on latest trades + const currentMarket = await Market.findOne({ marketId }); + let yesPrice = currentMarket ? currentMarket.currentYesPrice : 0.5; + let noPrice = currentMarket ? currentMarket.currentNoPrice : 0.5; + if (trade.tokenType === 'yes') { - update.currentYesPrice = trade.price; + yesPrice = trade.price; + noPrice = 1.0 - trade.price; } else if (trade.tokenType === 'no') { - update.currentNoPrice = trade.price; + noPrice = trade.price; + yesPrice = 1.0 - trade.price; } - await Market.findOneAndUpdate({ marketId }, update); + await eventSourcingService.appendEvent(marketId, 'PRICES_UPDATED', { yesPrice, noPrice }, trade.userWalletAddress); } catch (error) { logger.error('Failed to update market stats:', error); } diff --git a/backend/src/services/eventSourcingService.js b/backend/src/services/eventSourcingService.js new file mode 100644 index 0000000..0498a0e --- /dev/null +++ b/backend/src/services/eventSourcingService.js @@ -0,0 +1,229 @@ +const { Market, MarketEvent, MarketSnapshot } = require('../models'); +const logger = require('../config/logger'); + +class EventSourcingService { + /** + * Appends a new event to the event store and projects it onto the read model. + * + * @param {string} marketId The ID of the market. + * @param {string} eventType The type of the event (e.g., 'MARKET_CREATED'). + * @param {object} payload The event payload. + * @param {string} actorAddress The address of the actor initiating the event. + * @returns {object} The appended event. + */ + async appendEvent(marketId, eventType, payload, actorAddress = null) { + try { + // Determine the next sequence number for this market + const lastEvent = await MarketEvent.findOne({ marketId }) + .sort({ sequenceNumber: -1 }) + .limit(1); + + const sequenceNumber = lastEvent ? lastEvent.sequenceNumber + 1 : 1; + + // Create and save the new event + const event = new MarketEvent({ + marketId, + eventType, + schemaVersion: 1, + payload, + actorAddress, + sequenceNumber + }); + + await event.save(); + logger.debug('Appended event', { eventType, marketId, sequenceNumber }); + + // Synchronous projection to update the read model + await this.projectEvent(event); + + return event; + } catch (error) { + logger.error('Failed to append event', { eventType, marketId, error: error.message }); + throw error; + } + } + + /** + * Projects a single event onto the Market read model. + */ + async projectEvent(event) { + const { marketId, eventType, payload } = event; + + try { + switch (eventType) { + case 'MARKET_CREATED': { + const newMarket = new Market(payload); + await newMarket.save(); + break; + } + case 'MARKET_UPDATED': { + await Market.findOneAndUpdate({ marketId }, { $set: payload }); + break; + } + case 'MARKET_RESOLVED': { + await Market.findOneAndUpdate({ marketId }, { + status: 'resolved', + resolvedOutcome: payload.outcome, + resolvedAt: payload.resolvedAt || new Date(), + resolvedBy: payload.resolvedBy, + resolutionTxHash: payload.resolutionTxHash + }); + break; + } + case 'PRICES_UPDATED': { + const market = await Market.findOne({ marketId }); + if (market) { + market.updatePrices(payload.yesPrice, payload.noPrice); + await market.save(); + } + break; + } + case 'TRADE_EXECUTED': { + const market = await Market.findOne({ marketId }); + if (market) { + market.addTrade(payload.amount); + await market.save(); + } + break; + } + default: + logger.debug(`No projection handler for event type: ${eventType}`); + } + } catch (error) { + logger.error('Failed to project event', { eventType, marketId, error: error.message }); + throw error; + } + } + + /** + * Rebuilds the Market state from the event store. + * + * @param {string} marketId The ID of the market to rebuild. + * @returns {object} The reconstructed market data. + */ + async replayMarket(marketId) { + let state = {}; + let lastSequence = 0; + + // Check for a recent snapshot + const snapshot = await MarketSnapshot.findOne({ marketId }); + if (snapshot) { + state = { ...snapshot.stateData }; + lastSequence = snapshot.lastEventSequenceNumber; + } + + // Fetch all events after the snapshot + const events = await MarketEvent.find({ + marketId, + sequenceNumber: { $gt: lastSequence } + }).sort({ sequenceNumber: 1 }); + + for (const event of events) { + state = this.applyEventToState(state, event); + } + + return state; + } + + /** + * Applies an event to a raw state object (used during replay). + */ + applyEventToState(state, event) { + const { eventType, payload } = event; + + switch (eventType) { + case 'MARKET_CREATED': + return { ...payload }; + case 'MARKET_UPDATED': + return { ...state, ...payload }; + case 'MARKET_RESOLVED': + return { + ...state, + status: 'resolved', + resolvedOutcome: payload.outcome, + resolvedAt: payload.resolvedAt || new Date(), + resolvedBy: payload.resolvedBy, + resolutionTxHash: payload.resolutionTxHash + }; + case 'PRICES_UPDATED': + // NOTE: In a true pure function replay, you'd manage priceHistory array manually here + return { + ...state, + currentYesPrice: payload.yesPrice, + currentNoPrice: payload.noPrice, + }; + case 'TRADE_EXECUTED': + return { + ...state, + totalTrades: (state.totalTrades || 0) + 1, + totalVolume: (state.totalVolume || 0) + payload.amount, + }; + default: + return state; + } + } + + /** + * Generates a snapshot for a market to speed up future replays. + */ + async generateSnapshot(marketId) { + try { + const state = await this.replayMarket(marketId); + const lastEvent = await MarketEvent.findOne({ marketId }).sort({ sequenceNumber: -1 }).limit(1); + + if (!lastEvent) { + return null; + } + + const snapshot = await MarketSnapshot.findOneAndUpdate( + { marketId }, + { + stateData: state, + lastEventSequenceNumber: lastEvent.sequenceNumber, + timestamp: new Date() + }, + { upsert: true, new: true } + ); + + logger.info('Generated snapshot', { marketId, sequenceNumber: lastEvent.sequenceNumber }); + return snapshot; + } catch (error) { + logger.error('Failed to generate snapshot', { marketId, error: error.message }); + throw error; + } + } + + /** + * Recovers all markets by wiping the Market read models and rebuilding them from events. + * This is a dangerous operation typically only used in emergency recovery scenarios. + */ + async recoverAllMarkets() { + logger.warn('Starting full market recovery from event store...'); + try { + // Get all unique market IDs from the event store + const marketIds = await MarketEvent.distinct('marketId'); + + // Wipe current read models + await Market.deleteMany({}); + logger.info(`Cleared all existing market read models.`); + + let recoveredCount = 0; + for (const marketId of marketIds) { + const rebuiltState = await this.replayMarket(marketId); + if (rebuiltState && rebuiltState.marketId) { + const market = new Market(rebuiltState); + await market.save(); + recoveredCount++; + } + } + + logger.info(`Successfully recovered ${recoveredCount} markets from event store.`); + return { success: true, recoveredCount }; + } catch (error) { + logger.error('Market recovery failed:', error); + throw error; + } + } +} + +module.exports = new EventSourcingService();