diff --git a/backend/src/controllers/stream.controller.ts b/backend/src/controllers/stream.controller.ts index 60ac9301..f8888ac0 100644 --- a/backend/src/controllers/stream.controller.ts +++ b/backend/src/controllers/stream.controller.ts @@ -374,7 +374,7 @@ export const getStreamClaimableAmount = async (req: Request, res: Response) => { /** * Get user-level stream summary used by dashboard/profile cards. */ -export const getUserStreamSummary = async (req: Request, res: Response) => { +export const getUserStreamSummary = async (req: Request<{ address: string }>, res: Response) => { try { const address = Array.isArray(req.params.address) ? req.params.address[0] : (req.params.address ?? '').trim(); if (!address) { diff --git a/backend/src/services/sse.service.ts b/backend/src/services/sse.service.ts index f932da62..ec959484 100644 --- a/backend/src/services/sse.service.ts +++ b/backend/src/services/sse.service.ts @@ -20,7 +20,7 @@ interface SSECapacityCheckResult { message?: string; } -class SSEService { +export class SSEService { private clients: Map = new Map(); private readonly ipConnectionCounts: Map = new Map(); private shuttingDown = false; @@ -196,6 +196,13 @@ class SSEService { } } + broadcastToAdmin(event: string, data: unknown): void { + const adminKey = process.env.ADMIN_PUBLIC_KEY; + if (adminKey) { + this.broadcastToUser(adminKey, event, data); + } + } + private _localBroadcastToStream(streamId: string, event: string, data: unknown): void { this.broadcast(event, data, (client) => client.subscriptions.has(streamId) || client.subscriptions.has('*') @@ -225,5 +232,4 @@ class SSEService { } } -export { SSEService }; export const sseService = new SSEService(); diff --git a/backend/src/test/sseService.test.ts b/backend/src/test/sseService.test.ts deleted file mode 100644 index 01aed8ee..00000000 --- a/backend/src/test/sseService.test.ts +++ /dev/null @@ -1,87 +0,0 @@ -import { SseService } from "../services/sseService"; - -type MockRes = { - write: jest.Mock; - end: jest.Mock; -}; - -const createMockRes = (): MockRes => ({ - write: jest.fn(), - end: jest.fn(), -}); - -describe("SseService", () => { - let sse: SseService; - - beforeEach(() => { - sse = new SseService(); - }); - - test("test_subscribe_to_stream_events", () => { - const res = createMockRes(); - sse.subscribeToStream("stream-1", res as any); - - sse.broadcastToStream("stream-1", { msg: "hello" }); - - expect(res.write).toHaveBeenCalledWith( - expect.stringContaining("hello") - ); - }); - - test("test_subscribe_to_user_events", () => { - const res = createMockRes(); - sse.subscribeToUser("user-1", res as any); - - sse.broadcastToUser("user-1", { msg: "user event" }); - - expect(res.write).toHaveBeenCalledWith( - expect.stringContaining("user event") - ); - }); - - test("test_subscribe_all", () => { - const res = createMockRes(); - sse.subscribeAll(res as any); - - sse.broadcastAll({ msg: "global" }); - - expect(res.write).toHaveBeenCalledWith( - expect.stringContaining("global") - ); - }); - - test("test_client_disconnect_cleaned_up", () => { - const res = createMockRes(); - sse.subscribeToUser("user-1", res as any); - - sse.disconnect(res as any); - - expect(sse.getClientCount()).toBe(0); - }); - - test("test_broadcast_to_multiple_clients", () => { - const res1 = createMockRes(); - const res2 = createMockRes(); - - sse.subscribeToStream("stream-1", res1 as any); - sse.subscribeToStream("stream-1", res2 as any); - - sse.broadcastToStream("stream-1", { msg: "multi" }); - - expect(res1.write).toHaveBeenCalled(); - expect(res2.write).toHaveBeenCalled(); - }); - - test("test_no_cross_user_leakage", () => { - const resA = createMockRes(); - const resB = createMockRes(); - - sse.subscribeToUser("user-A", resA as any); - sse.subscribeToUser("user-B", resB as any); - - sse.broadcastToUser("user-A", { msg: "secret" }); - - expect(resA.write).toHaveBeenCalled(); - expect(resB.write).not.toHaveBeenCalled(); - }); -}); \ No newline at end of file diff --git a/backend/src/workers/soroban-event-worker.ts b/backend/src/workers/soroban-event-worker.ts index 29da161e..94788e3a 100644 --- a/backend/src/workers/soroban-event-worker.ts +++ b/backend/src/workers/soroban-event-worker.ts @@ -193,7 +193,18 @@ export class SorobanEventWorker { let lastCursor: string | null = state.lastCursor; let lastLedger: number = state.lastLedger; - for (const event of response.events) { + // Sort events so that 'stream_created' events are processed first in the batch. + // This ensures that subsequent events (like 'fee_collected') that depend on + // the stream existing in the DB can find it. + const sortedEvents = [...response.events].sort((a, b) => { + const aType = a.topic[0] ? decodeSymbol(a.topic[0]) : ''; + const bType = b.topic[0] ? decodeSymbol(b.topic[0]) : ''; + if (aType === 'stream_created' && bType !== 'stream_created') return -1; + if (bType === 'stream_created' && aType !== 'stream_created') return 1; + return 0; + }); + + for (const event of sortedEvents) { // Only process events from successful contract calls. if (!event.inSuccessfulContractCall) continue; @@ -607,7 +618,7 @@ export class SorobanEventWorker { }); // Broadcast to admin channel for treasury reporting - sseService.broadcast('stream.fee_collected', { + sseService.broadcastToAdmin('stream.fee_collected', { streamId, treasury, feeAmount, @@ -615,7 +626,7 @@ export class SorobanEventWorker { transactionHash: event.txHash, ledger: event.ledger, timestamp, - }, (client) => client.subscriptions.has('admin') || client.subscriptions.has('*')); + }); } private async handleStreamPaused( diff --git a/backend/tests/integration/streamInter.test.ts b/backend/tests/integration/streamInter.test.ts deleted file mode 100644 index dd3ed5f0..00000000 --- a/backend/tests/integration/streamInter.test.ts +++ /dev/null @@ -1,101 +0,0 @@ -import request from "supertest"; -import { app } from "../../../test/setup"; -import { db } from "../../db/client"; - -describe("Streams Integration", () => { - let streamId: string; - let sseMessages: any[] = []; - - // 🔌 Mock SSE client - const mockSseClient = () => { - return { - write: (data: string) => { - try { - const parsed = JSON.parse(data.replace(/^data:\s*/, "")); - sseMessages.push(parsed); - } catch {} - }, - end: jest.fn(), - }; - }; - - beforeEach(() => { - sseMessages = []; - }); - - test("POST /v1/streams creates stream + broadcasts SSE", async () => { - const sseClient = mockSseClient(); - app.get("sseService").subscribeAll(sseClient); - - const res = await request(app) - .post("/v1/streams") - .send({ - sender: "addr1", - recipient: "addr2", - amount: 1000, - rate: 1, - }); - - expect(res.status).toBe(201); - streamId = res.body.id; - - // ✅ DB check - const stream = await db.stream.findUnique({ where: { id: streamId } }); - expect(stream).toBeTruthy(); - - // ✅ SSE broadcast check - expect(sseMessages.length).toBeGreaterThan(0); - }); - - test("GET /v1/streams/{id} returns correct data", async () => { - const res = await request(app).get(`/v1/streams/${streamId}`); - - expect(res.status).toBe(200); - expect(res.body.id).toBe(streamId); - }); - - test("GET /v1/streams?sender filters correctly", async () => { - const res = await request(app) - .get("/v1/streams") - .query({ sender: "addr1" }); - - expect(res.status).toBe(200); - expect(res.body.every((s: any) => s.sender === "addr1")).toBe(true); - }); - - test("GET /v1/streams/{id}/events paginates", async () => { - const res = await request(app) - .get(`/v1/streams/${streamId}/events`) - .query({ limit: 10 }); - - expect(res.status).toBe(200); - expect(Array.isArray(res.body.data)).toBe(true); - }); - - test("Indexer processes TOPPED_UP event", async () => { - await request(app) - .post(`/v1/indexer/event`) - .send({ - type: "TOPPED_UP", - streamId, - amount: 500, - }); - - const stream = await db.stream.findUnique({ where: { id: streamId } }); - - expect(stream.depositedAmount).toBeGreaterThanOrEqual(1500); - }); - - test("Indexer processes CANCELLED event", async () => { - await request(app) - .post(`/v1/indexer/event`) - .send({ - type: "CANCELLED", - streamId, - }); - - const stream = await db.stream.findUnique({ where: { id: streamId } }); - - expect(stream.status).toBe("cancelled"); - }); -}); \ No newline at end of file