From d407fd1cc8c9188b313ce27d97bab8ef59d2a4a7 Mon Sep 17 00:00:00 2001 From: RedheadedProgrammer Date: Tue, 28 Jul 2026 20:35:43 +0100 Subject: [PATCH 1/5] feat(websockets): add Socket.IO server and horizon listener forwarding (fix #430) --- stellar-payment-platform/horizonListener.js | 28 ++++++++++- stellar-payment-platform/package.json | 4 +- stellar-payment-platform/server.js | 16 ++++++- stellar-payment-platform/src/socketManager.js | 47 +++++++++++++++++++ 4 files changed, 92 insertions(+), 3 deletions(-) create mode 100644 stellar-payment-platform/src/socketManager.js diff --git a/stellar-payment-platform/horizonListener.js b/stellar-payment-platform/horizonListener.js index 65b215d..66e4344 100644 --- a/stellar-payment-platform/horizonListener.js +++ b/stellar-payment-platform/horizonListener.js @@ -32,6 +32,20 @@ const POLL_INTERVAL_MS = parseInt(process.env.POLL_INTERVAL_MS, 10) || 60000; // // --------------------------------------------------------------------------- const horizon = new Horizon.Server(HORIZON_URL); +// Attempt to connect to the API server's Socket.IO endpoint so detected +// payments can be forwarded to connected clients in real-time. The listener +// runs as a separate process; it connects as a Socket.IO client and emits +// 'payment' events which the server will route to the appropriate room. +try { + const ioClient = require('socket.io-client'); + const SOCKET_SERVER = process.env.SOCKET_SERVER_URL || `http://localhost:${process.env.PORT || 5000}`; + global.__socketClient = ioClient(SOCKET_SERVER, { reconnection: true }); + global.__socketClient.on('connect', () => logger.info('[SocketClient] connected to server')); + global.__socketClient.on('connect_error', (err) => logger.error('[SocketClient] connect_error', err)); +} catch (err) { + logger.warn('Socket.IO client not available; real-time notifications disabled', err?.message || err); +} + // Track active streams so we can clean up on shutdown const activeStreams = new Map(); @@ -83,7 +97,19 @@ const watchAccount = (accountId) => { // Only log payment operations (ignore account_merge, etc.) if (payment.type === 'payment' || payment.type_i === 1) { logger.info(formatPayment(payment, accountId)); - } + + // If a Socket.IO server is available, emit the payment event so + // connected clients listening for this account receive a real-time + // notification. The horizon listener connects as a client to the + // API server and emits a 'payment' event with the address and payload. + try { + if (typeof global.__socketClient !== 'undefined' && global.__socketClient && global.__socketClient.connected) { + global.__socketClient.emit('payment', { address: accountId, payment }); + } + } catch (err) { + logger.error('Failed to emit payment over socket client', err); + } + } }, onerror: (error) => { logger.error( diff --git a/stellar-payment-platform/package.json b/stellar-payment-platform/package.json index fbcf619..2cb2b91 100644 --- a/stellar-payment-platform/package.json +++ b/stellar-payment-platform/package.json @@ -45,7 +45,9 @@ "winston": "^3.19.0", "winston-daily-rotate-file": "^5.0.0", "xss": "^1.0.15", - "zod": "^4.4.3" + "zod": "^4.4.3", + "socket.io": "^4.8.0", + "socket.io-client": "^4.8.0" }, "devDependencies": { "jest": "^29.7.0", diff --git a/stellar-payment-platform/server.js b/stellar-payment-platform/server.js index 3dd2afc..e262363 100644 --- a/stellar-payment-platform/server.js +++ b/stellar-payment-platform/server.js @@ -893,7 +893,21 @@ const gracefulShutdown = (server, prismaClient, signal) => { if (require.main === module) { - const server = app.listen(PORT, '0.0.0.0', () => { + const http = require('http'); + const { initSocketServer } = require('./src/socketManager'); + + const server = http.createServer(app); + + // Initialize Socket.IO on the HTTP server so the same port serves both + // the API and WebSocket connections. + try { + initSocketServer(server); + logger.info('Socket.IO initialized'); + } catch (err) { + logger.error('Failed to initialize Socket.IO', err); + } + + server.listen(PORT, '0.0.0.0', () => { logger.info(`Server successfully initialized on port ${PORT}`); }); diff --git a/stellar-payment-platform/src/socketManager.js b/stellar-payment-platform/src/socketManager.js new file mode 100644 index 0000000..e5dee29 --- /dev/null +++ b/stellar-payment-platform/src/socketManager.js @@ -0,0 +1,47 @@ +const { Server } = require('socket.io'); +const { logger } = require('./src/logger'); + +let io = null; + +/** + * Initialize Socket.IO on the passed http.Server instance. + * Maintains a simple room-per-address mapping: clients call `authenticate` with + * { address } and are added to a room named by that address. Emitted payments + * are broadcast to that room. + */ +function initSocketServer(httpServer, corsOptions = {}) { + if (io) return io; + io = new Server(httpServer, { + cors: Object.assign({ origin: true, methods: ['GET', 'POST'] }, corsOptions), + }); + + io.on('connection', (socket) => { + logger.info(`Socket connected: ${socket.id}`); + + socket.on('authenticate', (payload) => { + try { + const address = payload && typeof payload.address === 'string' ? payload.address : null; + if (address) { + socket.join(address); + socket.address = address; + logger.info(`Socket ${socket.id} joined room for ${address}`); + } + } catch (err) { + logger.error('Socket authenticate error', err); + } + }); + + socket.on('disconnect', () => { + logger.info(`Socket disconnected: ${socket.id}`); + }); + }); + + return io; +} + +function emitToAddress(address, event, data) { + if (!io) return; + io.to(address).emit(event, data); +} + +module.exports = { initSocketServer, emitToAddress }; From d63a39200bdb52c47e8affc698c99b0be4163b2f Mon Sep 17 00:00:00 2001 From: RedheadedProgrammer Date: Tue, 28 Jul 2026 20:45:58 +0100 Subject: [PATCH 2/5] fix(websockets): correct logger require path in socketManager.js --- stellar-payment-platform/src/socketManager.js | 2 +- 1 file changed, 1 insertion(+), 1 deletion(-) diff --git a/stellar-payment-platform/src/socketManager.js b/stellar-payment-platform/src/socketManager.js index e5dee29..316464f 100644 --- a/stellar-payment-platform/src/socketManager.js +++ b/stellar-payment-platform/src/socketManager.js @@ -1,5 +1,5 @@ const { Server } = require('socket.io'); -const { logger } = require('./src/logger'); +const { logger } = require('./logger'); let io = null; From bcabc0a0c6d60db053b2492096c8981cc53e0d2e Mon Sep 17 00:00:00 2001 From: RedheadedProgrammer Date: Thu, 30 Jul 2026 13:11:14 +0100 Subject: [PATCH 3/5] fix(websockets): handle payment forwarding event and update lockfile --- stellar-payment-platform/package-lock.json | 303 +++++++++++++++++- stellar-payment-platform/src/socketManager.js | 21 +- stellar-payment-platform/tests/socket.test.js | 70 ++++ 3 files changed, 391 insertions(+), 3 deletions(-) create mode 100644 stellar-payment-platform/tests/socket.test.js diff --git a/stellar-payment-platform/package-lock.json b/stellar-payment-platform/package-lock.json index 62f2312..57d08ca 100644 --- a/stellar-payment-platform/package-lock.json +++ b/stellar-payment-platform/package-lock.json @@ -28,6 +28,8 @@ "prom-client": "^15.1.3", "rate-limit-redis": "^4.2.0", "redis": "^4.7.0", + "socket.io": "^4.8.0", + "socket.io-client": "^4.8.0", "sqlite3": "^5.1.7", "uuid": "^9.0.1", "winston": "^3.19.0", @@ -2004,6 +2006,12 @@ "text-hex": "1.0.x" } }, + "node_modules/@socket.io/component-emitter": { + "version": "3.1.2", + "resolved": "https://registry.npmjs.org/@socket.io/component-emitter/-/component-emitter-3.1.2.tgz", + "integrity": "sha512-9BCxFwvbGg/RsZK9tjXd8s4UcwR0MWeFQ1XEKIQVVvAGJyINdrqKMcTRyLoK8Rse1GjzLV9cwjWV1olXRWEXVA==", + "license": "MIT" + }, "node_modules/@standard-schema/spec": { "version": "1.1.0", "resolved": "https://registry.npmjs.org/@standard-schema/spec/-/spec-1.1.0.tgz", @@ -2111,6 +2119,15 @@ "@babel/types": "^7.28.2" } }, + "node_modules/@types/cors": { + "version": "2.8.19", + "resolved": "https://registry.npmjs.org/@types/cors/-/cors-2.8.19.tgz", + "integrity": "sha512-mFNylyeyqN93lfe/9CSxOGREz8cpzAhH+E93xJ4xWQf62V8sQ/24reV2nyzUWM6H6Xji+GGHpkbLe7pVoUEskg==", + "license": "MIT", + "dependencies": { + "@types/node": "*" + } + }, "node_modules/@types/estree": { "version": "1.0.9", "resolved": "https://registry.npmjs.org/@types/estree/-/estree-1.0.9.tgz", @@ -2158,7 +2175,6 @@ "version": "26.0.1", "resolved": "https://registry.npmjs.org/@types/node/-/node-26.0.1.tgz", "integrity": "sha512-fc3KiUoBt6kie0N9bIW3E47vZsuaMf0PM2AaUpLCLT0s/LvX1nxAim6Fc049cNxODPpGm6qRAuUOB86SkRuPQw==", - "dev": true, "license": "MIT", "dependencies": { "undici-types": "~8.3.0" @@ -2177,6 +2193,15 @@ "integrity": "sha512-6WaYesThRMCl19iryMYP7/x2OVgCtbIVflDGFpWnb9irXI3UjYE4AzmYuiUKY1AJstGijoY+MgUszMgRxIYTYw==", "license": "MIT" }, + "node_modules/@types/ws": { + "version": "8.18.1", + "resolved": "https://registry.npmjs.org/@types/ws/-/ws-8.18.1.tgz", + "integrity": "sha512-ThVF6DCVhA8kUGy+aazFQ4kXQ7E1Ty7A3ypFOe0IcJV8O/M511G99AW24irKrW56Wt44yG9+ij8FaqoBGkuBXg==", + "license": "MIT", + "dependencies": { + "@types/node": "*" + } + }, "node_modules/@types/yargs": { "version": "17.0.35", "resolved": "https://registry.npmjs.org/@types/yargs/-/yargs-17.0.35.tgz", @@ -2620,6 +2645,15 @@ ], "license": "MIT" }, + "node_modules/base64id": { + "version": "2.0.0", + "resolved": "https://registry.npmjs.org/base64id/-/base64id-2.0.0.tgz", + "integrity": "sha512-lGe34o6EHj9y3Kts9R4ZYs/Gr+6N7MCaMlIFA3F1R2O5/m7K06AxfSeO5530PEERE6/WyEg3lsuyw4GHlPZHog==", + "license": "MIT", + "engines": { + "node": "^4.5.0 || >= 5.9" + } + }, "node_modules/baseline-browser-mapping": { "version": "2.10.40", "resolved": "https://registry.npmjs.org/baseline-browser-mapping/-/baseline-browser-mapping-2.10.40.tgz", @@ -3895,6 +3929,95 @@ "once": "^1.4.0" } }, + "node_modules/engine.io": { + "version": "6.6.9", + "resolved": "https://registry.npmjs.org/engine.io/-/engine.io-6.6.9.tgz", + "integrity": "sha512-clKkw4C7nJ22mGgoVcCg6V/W/TxdNyIOTr89k2ONZu81qqkddPFDF0LXcbAwhzPD8DjkiRCjzuiO6Y+fkpD4vg==", + "license": "MIT", + "dependencies": { + "@types/cors": "^2.8.12", + "@types/node": ">=10.0.0", + "@types/ws": "^8.5.12", + "accepts": "~1.3.4", + "base64id": "2.0.0", + "cookie": "~0.7.2", + "cors": "~2.8.5", + "debug": "~4.4.1", + "engine.io-parser": "~5.2.1", + "ws": "~8.21.0" + }, + "engines": { + "node": ">=10.2.0" + } + }, + "node_modules/engine.io-client": { + "version": "6.6.6", + "resolved": "https://registry.npmjs.org/engine.io-client/-/engine.io-client-6.6.6.tgz", + "integrity": "sha512-iY6QdftLQ9pyiPoX082bpf/u1UewnOaJrtJIF9T0++QB34lZrj0uP+Q/bj8AlUsAxqhnkTV2BS8SBZSxOmoV5Q==", + "license": "MIT", + "dependencies": { + "@socket.io/component-emitter": "~3.1.0", + "debug": "~4.4.1", + "engine.io-parser": "~5.2.1", + "ws": "~8.21.0", + "xmlhttprequest-ssl": "~2.1.1" + } + }, + "node_modules/engine.io-client/node_modules/debug": { + "version": "4.4.3", + "resolved": "https://registry.npmjs.org/debug/-/debug-4.4.3.tgz", + "integrity": "sha512-RGwwWnwQvkVfavKVt22FGLw+xYSdzARwm0ru6DhTVA3umU5hZc28V3kO4stgYryrTlLpuvgI9GiijltAjNbcqA==", + "license": "MIT", + "dependencies": { + "ms": "^2.1.3" + }, + "engines": { + "node": ">=6.0" + }, + "peerDependenciesMeta": { + "supports-color": { + "optional": true + } + } + }, + "node_modules/engine.io-client/node_modules/ms": { + "version": "2.1.3", + "resolved": "https://registry.npmjs.org/ms/-/ms-2.1.3.tgz", + "integrity": "sha512-6FlzubTLZG3J2a/NVCAleEhjzq5oxgHyaCU9yYXvcLsvoVaHJq/s5xXI6/XXP6tz7R9xAOtHnSO/tXtF3WRTlA==", + "license": "MIT" + }, + "node_modules/engine.io-parser": { + "version": "5.2.3", + "resolved": "https://registry.npmjs.org/engine.io-parser/-/engine.io-parser-5.2.3.tgz", + "integrity": "sha512-HqD3yTBfnBxIrbnM1DoD6Pcq8NECnh8d4As1Qgh0z5Gg3jRRIqijury0CL3ghu/edArpUYiYqQiDUQBIs4np3Q==", + "license": "MIT", + "engines": { + "node": ">=10.0.0" + } + }, + "node_modules/engine.io/node_modules/debug": { + "version": "4.4.3", + "resolved": "https://registry.npmjs.org/debug/-/debug-4.4.3.tgz", + "integrity": "sha512-RGwwWnwQvkVfavKVt22FGLw+xYSdzARwm0ru6DhTVA3umU5hZc28V3kO4stgYryrTlLpuvgI9GiijltAjNbcqA==", + "license": "MIT", + "dependencies": { + "ms": "^2.1.3" + }, + "engines": { + "node": ">=6.0" + }, + "peerDependenciesMeta": { + "supports-color": { + "optional": true + } + } + }, + "node_modules/engine.io/node_modules/ms": { + "version": "2.1.3", + "resolved": "https://registry.npmjs.org/ms/-/ms-2.1.3.tgz", + "integrity": "sha512-6FlzubTLZG3J2a/NVCAleEhjzq5oxgHyaCU9yYXvcLsvoVaHJq/s5xXI6/XXP6tz7R9xAOtHnSO/tXtF3WRTlA==", + "license": "MIT" + }, "node_modules/env-paths": { "version": "2.2.1", "resolved": "https://registry.npmjs.org/env-paths/-/env-paths-2.2.1.tgz", @@ -8165,6 +8288,154 @@ "url": "https://github.com/sponsors/cyyynthia" } }, + "node_modules/socket.io": { + "version": "4.8.3", + "resolved": "https://registry.npmjs.org/socket.io/-/socket.io-4.8.3.tgz", + "integrity": "sha512-2Dd78bqzzjE6KPkD5fHZmDAKRNe3J15q+YHDrIsy9WEkqttc7GY+kT9OBLSMaPbQaEd0x1BjcmtMtXkfpc+T5A==", + "license": "MIT", + "dependencies": { + "accepts": "~1.3.4", + "base64id": "~2.0.0", + "cors": "~2.8.5", + "debug": "~4.4.1", + "engine.io": "~6.6.0", + "socket.io-adapter": "~2.5.2", + "socket.io-parser": "~4.2.4" + }, + "engines": { + "node": ">=10.2.0" + } + }, + "node_modules/socket.io-adapter": { + "version": "2.5.8", + "resolved": "https://registry.npmjs.org/socket.io-adapter/-/socket.io-adapter-2.5.8.tgz", + "integrity": "sha512-6Oy52pbg+kvdCVvjcN+FnY7BvxZ7cIHNScbvztT/It5d0vbwoJoVZmF2gjJmnV0/4WlXRfG15zc45ySk9Ah8bw==", + "license": "MIT", + "dependencies": { + "debug": "~4.4.1", + "ws": "~8.21.0" + } + }, + "node_modules/socket.io-adapter/node_modules/debug": { + "version": "4.4.3", + "resolved": "https://registry.npmjs.org/debug/-/debug-4.4.3.tgz", + "integrity": "sha512-RGwwWnwQvkVfavKVt22FGLw+xYSdzARwm0ru6DhTVA3umU5hZc28V3kO4stgYryrTlLpuvgI9GiijltAjNbcqA==", + "license": "MIT", + "dependencies": { + "ms": "^2.1.3" + }, + "engines": { + "node": ">=6.0" + }, + "peerDependenciesMeta": { + "supports-color": { + "optional": true + } + } + }, + "node_modules/socket.io-adapter/node_modules/ms": { + "version": "2.1.3", + "resolved": "https://registry.npmjs.org/ms/-/ms-2.1.3.tgz", + "integrity": "sha512-6FlzubTLZG3J2a/NVCAleEhjzq5oxgHyaCU9yYXvcLsvoVaHJq/s5xXI6/XXP6tz7R9xAOtHnSO/tXtF3WRTlA==", + "license": "MIT" + }, + "node_modules/socket.io-client": { + "version": "4.8.3", + "resolved": "https://registry.npmjs.org/socket.io-client/-/socket.io-client-4.8.3.tgz", + "integrity": "sha512-uP0bpjWrjQmUt5DTHq9RuoCBdFJF10cdX9X+a368j/Ft0wmaVgxlrjvK3kjvgCODOMMOz9lcaRzxmso0bTWZ/g==", + "license": "MIT", + "dependencies": { + "@socket.io/component-emitter": "~3.1.0", + "debug": "~4.4.1", + "engine.io-client": "~6.6.1", + "socket.io-parser": "~4.2.4" + }, + "engines": { + "node": ">=10.0.0" + } + }, + "node_modules/socket.io-client/node_modules/debug": { + "version": "4.4.3", + "resolved": "https://registry.npmjs.org/debug/-/debug-4.4.3.tgz", + "integrity": "sha512-RGwwWnwQvkVfavKVt22FGLw+xYSdzARwm0ru6DhTVA3umU5hZc28V3kO4stgYryrTlLpuvgI9GiijltAjNbcqA==", + "license": "MIT", + "dependencies": { + "ms": "^2.1.3" + }, + "engines": { + "node": ">=6.0" + }, + "peerDependenciesMeta": { + "supports-color": { + "optional": true + } + } + }, + "node_modules/socket.io-client/node_modules/ms": { + "version": "2.1.3", + "resolved": "https://registry.npmjs.org/ms/-/ms-2.1.3.tgz", + "integrity": "sha512-6FlzubTLZG3J2a/NVCAleEhjzq5oxgHyaCU9yYXvcLsvoVaHJq/s5xXI6/XXP6tz7R9xAOtHnSO/tXtF3WRTlA==", + "license": "MIT" + }, + "node_modules/socket.io-parser": { + "version": "4.2.7", + "resolved": "https://registry.npmjs.org/socket.io-parser/-/socket.io-parser-4.2.7.tgz", + "integrity": "sha512-IH/iSeO9T6gz1KkFleGDWkG9N3dl4jXVYUtMhIqH10Md0ttMer8nUNWiP1DKuNrybD2xBrixLJdCC9J6ECoYkg==", + "license": "MIT", + "dependencies": { + "@socket.io/component-emitter": "~3.1.0", + "debug": "~4.4.1" + }, + "engines": { + "node": ">=10.0.0" + } + }, + "node_modules/socket.io-parser/node_modules/debug": { + "version": "4.4.3", + "resolved": "https://registry.npmjs.org/debug/-/debug-4.4.3.tgz", + "integrity": "sha512-RGwwWnwQvkVfavKVt22FGLw+xYSdzARwm0ru6DhTVA3umU5hZc28V3kO4stgYryrTlLpuvgI9GiijltAjNbcqA==", + "license": "MIT", + "dependencies": { + "ms": "^2.1.3" + }, + "engines": { + "node": ">=6.0" + }, + "peerDependenciesMeta": { + "supports-color": { + "optional": true + } + } + }, + "node_modules/socket.io-parser/node_modules/ms": { + "version": "2.1.3", + "resolved": "https://registry.npmjs.org/ms/-/ms-2.1.3.tgz", + "integrity": "sha512-6FlzubTLZG3J2a/NVCAleEhjzq5oxgHyaCU9yYXvcLsvoVaHJq/s5xXI6/XXP6tz7R9xAOtHnSO/tXtF3WRTlA==", + "license": "MIT" + }, + "node_modules/socket.io/node_modules/debug": { + "version": "4.4.3", + "resolved": "https://registry.npmjs.org/debug/-/debug-4.4.3.tgz", + "integrity": "sha512-RGwwWnwQvkVfavKVt22FGLw+xYSdzARwm0ru6DhTVA3umU5hZc28V3kO4stgYryrTlLpuvgI9GiijltAjNbcqA==", + "license": "MIT", + "dependencies": { + "ms": "^2.1.3" + }, + "engines": { + "node": ">=6.0" + }, + "peerDependenciesMeta": { + "supports-color": { + "optional": true + } + } + }, + "node_modules/socket.io/node_modules/ms": { + "version": "2.1.3", + "resolved": "https://registry.npmjs.org/ms/-/ms-2.1.3.tgz", + "integrity": "sha512-6FlzubTLZG3J2a/NVCAleEhjzq5oxgHyaCU9yYXvcLsvoVaHJq/s5xXI6/XXP6tz7R9xAOtHnSO/tXtF3WRTlA==", + "license": "MIT" + }, "node_modules/socks": { "version": "2.8.9", "resolved": "https://registry.npmjs.org/socks/-/socks-2.8.9.tgz", @@ -8789,7 +9060,6 @@ "version": "8.3.0", "resolved": "https://registry.npmjs.org/undici-types/-/undici-types-8.3.0.tgz", "integrity": "sha512-j375ScV60dom+YkPFIfTLcOiPxkN/buHz5GobjLhixFuANaNs3C9l4GmrWqejgXWJ7BbJcFYpTEUkS1Ge8bpZQ==", - "dev": true, "license": "MIT" }, "node_modules/unicode-properties": { @@ -9117,6 +9387,35 @@ "node": "^12.13.0 || ^14.15.0 || >=16.0.0" } }, + "node_modules/ws": { + "version": "8.21.1", + "resolved": "https://registry.npmjs.org/ws/-/ws-8.21.1.tgz", + "integrity": "sha512-+0NTnW77fFN/DjQi6k/Sq/Yvk4Sgajw7urW8V+asjXnRgDs9gyGkdb7EzgfhA4goXsRIZKE28fzIXBHEzhuiWw==", + "license": "MIT", + "engines": { + "node": ">=10.0.0" + }, + "peerDependencies": { + "bufferutil": "^4.0.1", + "utf-8-validate": ">=5.0.2" + }, + "peerDependenciesMeta": { + "bufferutil": { + "optional": true + }, + "utf-8-validate": { + "optional": true + } + } + }, + "node_modules/xmlhttprequest-ssl": { + "version": "2.1.2", + "resolved": "https://registry.npmjs.org/xmlhttprequest-ssl/-/xmlhttprequest-ssl-2.1.2.tgz", + "integrity": "sha512-TEU+nJVUUnA4CYJFLvK5X9AOeH4KvDvhIfm0vV1GaQRtchnG0hgK5p8hw/xjv8cunWYCsiPCSDzObPyhEwq3KQ==", + "engines": { + "node": ">=0.4.0" + } + }, "node_modules/xss": { "version": "1.0.15", "resolved": "https://registry.npmjs.org/xss/-/xss-1.0.15.tgz", diff --git a/stellar-payment-platform/src/socketManager.js b/stellar-payment-platform/src/socketManager.js index 316464f..62490ba 100644 --- a/stellar-payment-platform/src/socketManager.js +++ b/stellar-payment-platform/src/socketManager.js @@ -31,6 +31,18 @@ function initSocketServer(httpServer, corsOptions = {}) { } }); + socket.on('payment', (payload) => { + try { + const { address, payment } = payload || {}; + if (address && payment) { + emitToAddress(address, 'payment', payment); + logger.info(`Forwarded payment event to room for address ${address}`); + } + } catch (err) { + logger.error('Socket payment forwarding error', err); + } + }); + socket.on('disconnect', () => { logger.info(`Socket disconnected: ${socket.id}`); }); @@ -44,4 +56,11 @@ function emitToAddress(address, event, data) { io.to(address).emit(event, data); } -module.exports = { initSocketServer, emitToAddress }; +function closeSocketServer() { + if (io) { + io.close(); + io = null; + } +} + +module.exports = { initSocketServer, emitToAddress, closeSocketServer }; diff --git a/stellar-payment-platform/tests/socket.test.js b/stellar-payment-platform/tests/socket.test.js new file mode 100644 index 0000000..fb27fe3 --- /dev/null +++ b/stellar-payment-platform/tests/socket.test.js @@ -0,0 +1,70 @@ +'use strict'; + +const http = require('http'); +const ioClient = require('socket.io-client'); +const { initSocketServer, closeSocketServer } = require('../src/socketManager'); + +// Set a larger timeout for the socket tests to avoid intermittent timeouts under heavy test run loads +jest.setTimeout(30000); + +// Mock logger to avoid spamming output +jest.mock('../src/logger', () => ({ + logger: { + info: jest.fn(), + error: jest.fn(), + warn: jest.fn(), + }, +})); + +describe('Socket.IO real-time notification system', () => { + let server; + let clientSocket; + let port; + + beforeAll((done) => { + server = http.createServer(); + initSocketServer(server); + server.listen(() => { + port = server.address().port; + done(); + }); + }); + + afterAll((done) => { + if (clientSocket && clientSocket.connected) { + clientSocket.disconnect(); + } + closeSocketServer(); + server.close(done); + }); + + beforeEach((done) => { + clientSocket = ioClient(`http://localhost:${port}`); + clientSocket.on('connect', done); + }); + + afterEach(() => { + clientSocket.disconnect(); + }); + + test('should authenticate and join the address room and receive payment events', (done) => { + const testAddress = 'GAPUQZH3WZUXHEMUGZN5ZYU4D4GHCFEMOGUINU6MF345GBD2QXNYYIEQ'; + clientSocket.emit('authenticate', { address: testAddress }); + + setTimeout(() => { + clientSocket.on('payment', (paymentData) => { + expect(paymentData).toEqual({ amount: '100', asset: 'XLM' }); + done(); + }); + + const publisherSocket = ioClient(`http://localhost:${port}`); + publisherSocket.on('connect', () => { + publisherSocket.emit('payment', { + address: testAddress, + payment: { amount: '100', asset: 'XLM' }, + }); + publisherSocket.disconnect(); + }); + }, 200); + }); +}); From b659b1b142d959a5608d85b0430aa815b71d5f31 Mon Sep 17 00:00:00 2001 From: RedheadedProgrammer Date: Mon, 3 Aug 2026 10:28:47 +0100 Subject: [PATCH 4/5] Merge websocket changes: resolve conflicts in package.json, horizonListener.js and socketManager.js --- stellar-payment-platform/horizonListener.js | 3 +++ stellar-payment-platform/package.json | 3 +++ stellar-payment-platform/src/socketManager.js | 23 +++++++++++++------ 3 files changed, 22 insertions(+), 7 deletions(-) diff --git a/stellar-payment-platform/horizonListener.js b/stellar-payment-platform/horizonListener.js index 3b1078e..dc8a5be 100644 --- a/stellar-payment-platform/horizonListener.js +++ b/stellar-payment-platform/horizonListener.js @@ -114,6 +114,7 @@ const watchAccount = (accountId) => { logger.error('Failed to emit payment over socket client', err); } } +<<<<<<< HEAD dispatchPaymentWebhooks({ prisma, poolGetFn: poolGet, @@ -126,6 +127,8 @@ const watchAccount = (accountId) => { ), ); } +======= +>>>>>>> 4095f79 (feat(websockets): add Socket.IO server and horizon listener forwarding (fix #430)) }, onerror: (error) => { logger.error( diff --git a/stellar-payment-platform/package.json b/stellar-payment-platform/package.json index 5363c6b..560d47a 100644 --- a/stellar-payment-platform/package.json +++ b/stellar-payment-platform/package.json @@ -45,7 +45,10 @@ "winston": "^3.19.0", "winston-daily-rotate-file": "^5.0.0", "xss": "^1.0.15", +<<<<<<< HEAD "zod": "^4.4.3", +======= +>>>>>>> 4095f79 (feat(websockets): add Socket.IO server and horizon listener forwarding (fix #430)) "socket.io": "^4.8.0", "socket.io-client": "^4.8.0" }, diff --git a/stellar-payment-platform/src/socketManager.js b/stellar-payment-platform/src/socketManager.js index 62490ba..21ed493 100644 --- a/stellar-payment-platform/src/socketManager.js +++ b/stellar-payment-platform/src/socketManager.js @@ -18,28 +18,37 @@ function initSocketServer(httpServer, corsOptions = {}) { io.on('connection', (socket) => { logger.info(`Socket connected: ${socket.id}`); - socket.on('authenticate', (payload) => { + // The optional ack lets a client wait until it is actually subscribed. + // Without it a client that emits `authenticate` and immediately expects + // events can miss any payment that arrives before the room join lands. + socket.on('authenticate', async (payload, ack) => { + let subscribed = false; try { const address = payload && typeof payload.address === 'string' ? payload.address : null; if (address) { - socket.join(address); + // socket.join can be async in some adapters; await to be safe. + await socket.join(address); socket.address = address; + subscribed = true; logger.info(`Socket ${socket.id} joined room for ${address}`); } } catch (err) { logger.error('Socket authenticate error', err); } + if (typeof ack === 'function') { + try { ack({ subscribed }); } catch (e) { /* ignore ack errors */ } + } }); socket.on('payment', (payload) => { try { - const { address, payment } = payload || {}; - if (address && payment) { - emitToAddress(address, 'payment', payment); - logger.info(`Forwarded payment event to room for address ${address}`); + const addr = payload && typeof payload.address === 'string' ? payload.address : null; + const payment = payload && payload.payment ? payload.payment : payload; + if (addr) { + io.to(addr).emit('payment', payment); } } catch (err) { - logger.error('Socket payment forwarding error', err); + logger.error('Error handling payment emit', err); } }); From 11584e7fc62a95b9599277c73e8967ed2a1bc7e7 Mon Sep 17 00:00:00 2001 From: RedheadedProgrammer Date: Mon, 3 Aug 2026 10:34:34 +0100 Subject: [PATCH 5/5] Resolve merge conflicts for websocket tests and socket manager --- stellar-payment-platform/src/socketManager.js | 2 +- stellar-payment-platform/tests/socket.test.js | 115 ++++++++++++------ 2 files changed, 78 insertions(+), 39 deletions(-) diff --git a/stellar-payment-platform/src/socketManager.js b/stellar-payment-platform/src/socketManager.js index 21ed493..26a9c65 100644 --- a/stellar-payment-platform/src/socketManager.js +++ b/stellar-payment-platform/src/socketManager.js @@ -1,4 +1,4 @@ -const { Server } = require('socket.io'); +const { Server } = require('socket.io'); const { logger } = require('./logger'); let io = null; diff --git a/stellar-payment-platform/tests/socket.test.js b/stellar-payment-platform/tests/socket.test.js index fb27fe3..42817b2 100644 --- a/stellar-payment-platform/tests/socket.test.js +++ b/stellar-payment-platform/tests/socket.test.js @@ -1,70 +1,109 @@ -'use strict'; +jest.setTimeout(20000); +jest.mock('../src/logger', () => ({ logger: { info: jest.fn(), warn: jest.fn(), error: jest.fn() } })); const http = require('http'); const ioClient = require('socket.io-client'); -const { initSocketServer, closeSocketServer } = require('../src/socketManager'); +const { initSocketServer } = require('../src/socketManager'); -// Set a larger timeout for the socket tests to avoid intermittent timeouts under heavy test run loads -jest.setTimeout(30000); +let server; +let clientSocket; +let publisherSocket; +let port; -// Mock logger to avoid spamming output -jest.mock('../src/logger', () => ({ - logger: { - info: jest.fn(), - error: jest.fn(), - warn: jest.fn(), - }, -})); +function startSocketServer(done) { + // Create a minimal express app for the socket server so we don't pull in + // heavyweight dependencies (like @stellar/stellar-sdk) during tests. + const express = require('express'); + const app = express(); + app.get('/health', (_req, res) => res.json({ ok: true })); -describe('Socket.IO real-time notification system', () => { - let server; - let clientSocket; - let port; + const httpServer = http.createServer(app); + initSocketServer(httpServer); + httpServer.listen(0, '127.0.0.1', () => { + port = httpServer.address().port; + server = httpServer; + done(); + }); +} + +function closeSocketServer() { + if (server && server.listening) server.close(); +} +describe('Socket.IO real-time notification system', () => { beforeAll((done) => { - server = http.createServer(); - initSocketServer(server); - server.listen(() => { - port = server.address().port; - done(); + startSocketServer(() => { + clientSocket = ioClient(`http://localhost:${port}`, { transports: ['websocket'] }); + clientSocket.on('connect', () => done()); }); }); afterAll((done) => { - if (clientSocket && clientSocket.connected) { - clientSocket.disconnect(); - } - closeSocketServer(); - server.close(done); - }); + // Ensure client sockets are disconnected before closing the server so + // server.close's callback fires promptly. + try { + if (clientSocket && clientSocket.connected) clientSocket.disconnect(); + if (publisherSocket && publisherSocket.connected) publisherSocket.disconnect(); + } catch (e) { /* ignore */ } - beforeEach((done) => { - clientSocket = ioClient(`http://localhost:${port}`); - clientSocket.on('connect', done); + if (server && server.listening) { + server.close(() => done()); + return; + } + done(); }); afterEach(() => { - clientSocket.disconnect(); + if (publisherSocket && publisherSocket.connected) publisherSocket.disconnect(); + publisherSocket = null; }); test('should authenticate and join the address room and receive payment events', (done) => { const testAddress = 'GAPUQZH3WZUXHEMUGZN5ZYU4D4GHCFEMOGUINU6MF345GBD2QXNYYIEQ'; - clientSocket.emit('authenticate', { address: testAddress }); - setTimeout(() => { - clientSocket.on('payment', (paymentData) => { + clientSocket.on('payment', (paymentData) => { + try { expect(paymentData).toEqual({ amount: '100', asset: 'XLM' }); done(); - }); + } catch (err) { + done(err); + } + }); - const publisherSocket = ioClient(`http://localhost:${port}`); + // The ack fires after the server has joined the room, so the publish below + // cannot outrun the subscription. + clientSocket.emit('authenticate', { address: testAddress }, () => { + publisherSocket = ioClient(`http://localhost:${port}`, { transports: ['websocket'] }); publisherSocket.on('connect', () => { publisherSocket.emit('payment', { address: testAddress, payment: { amount: '100', asset: 'XLM' }, }); - publisherSocket.disconnect(); + // Disconnecting here would close the transport before the packet is + // flushed and drop the event; afterEach tears the socket down instead. }); - }, 200); + }); + }); + + test('does not deliver payments for an address the client did not subscribe to', (done) => { + const subscribed = 'GAPUQZH3WZUXHEMUGZN5ZYU4D4GHCFEMOGUINU6MF345GBD2QXNYYIEQ'; + const other = 'GBDQD3WTQ6W2VQ2W4V74UZ5WYF6B72GZ6EHD7I3L3WYH357Y4K5H3E4W'; + + clientSocket.on('payment', () => { + done(new Error('received a payment addressed to another account')); + }); + + clientSocket.emit('authenticate', { address: subscribed }, () => { + publisherSocket = ioClient(`http://localhost:${port}`, { transports: ['websocket'] }); + publisherSocket.on('connect', () => { + publisherSocket.emit('payment', { + address: other, + payment: { amount: '100', asset: 'XLM' }, + }); + // Round-trip through the same socket: once this ack returns the server + // has already handled the payment above, so nothing is still in flight. + publisherSocket.emit('authenticate', { address: other }, () => done()); + }); + }); }); });