From a58994df3f7141b4ddd852fae99db8ae273822c2 Mon Sep 17 00:00:00 2001 From: onyekachi66 Date: Fri, 28 Aug 2026 12:40:41 +0100 Subject: [PATCH] feat: persist stream events in indexer --- indexer/streams/src/db/repository.ts | 133 ++++++++ indexer/streams/src/handlers/index.ts | 13 +- indexer/streams/src/handlers/persistence.ts | 29 ++ .../src/handlers/stream-cancel.handler.ts | 72 +++-- .../src/handlers/stream-created.handler.ts | 58 ++++ .../src/handlers/stream-funded.handler.ts | 83 ++--- .../src/handlers/stream-handlers.test.ts | 292 ++++++++++++------ .../src/handlers/stream-withdrawal.handler.ts | 71 +++-- .../src/handlers/streamCreated.test.ts | 179 ----------- indexer/streams/src/handlers/streamCreated.ts | 88 ------ indexer/streams/src/handlers/types.ts | 25 ++ indexer/streams/src/index.test.ts | 13 - indexer/streams/src/index.ts | 6 +- 13 files changed, 589 insertions(+), 473 deletions(-) create mode 100644 indexer/streams/src/db/repository.ts create mode 100644 indexer/streams/src/handlers/persistence.ts create mode 100644 indexer/streams/src/handlers/stream-created.handler.ts delete mode 100644 indexer/streams/src/handlers/streamCreated.test.ts delete mode 100644 indexer/streams/src/handlers/streamCreated.ts delete mode 100644 indexer/streams/src/index.test.ts diff --git a/indexer/streams/src/db/repository.ts b/indexer/streams/src/db/repository.ts new file mode 100644 index 0000000..fea9d38 --- /dev/null +++ b/indexer/streams/src/db/repository.ts @@ -0,0 +1,133 @@ +import type { DataSource } from "typeorm"; +import { Stream } from "./entity/Stream.js"; +import { WithdrawalAction } from "./entity/WithdrawalAction.js"; +import { CancelAction } from "./entity/CancelAction.js"; + +export interface CreateStreamInput { + id: string; + sender: string; + recipient: string; + token: string; + totalAmount: string; + startTime: string; + endTime: string; +} + +export interface FundStreamInput { + streamId: string; + amount: string; +} + +export interface RecordWithdrawalInput { + streamId: string; + recipient: string; + amount: string; + txHash: string; + timestamp: string; +} + +export interface RecordCancelInput { + streamId: string; + canceler: string; + txHash: string; + timestamp: string; +} + +export interface StreamPersistence { + createStream(input: CreateStreamInput): Promise; + fundStream(input: FundStreamInput): Promise; + recordWithdrawal(input: RecordWithdrawalInput): Promise; + recordCancel(input: RecordCancelInput): Promise; +} + +export class StreamRepository implements StreamPersistence { + constructor(private readonly dataSource: DataSource) {} + + async createStream(input: CreateStreamInput): Promise { + await this.dataSource + .getRepository(Stream) + .createQueryBuilder() + .insert() + .into(Stream) + .values({ + id: input.id, + sender: input.sender, + recipient: input.recipient, + token: input.token, + totalAmount: input.totalAmount, + startTime: input.startTime, + endTime: input.endTime, + amountWithdrawn: "0", + canceled: false, + }) + .orIgnore() // Ignore duplicates on primary key + .execute(); + } + + async fundStream(input: FundStreamInput): Promise { + await this.dataSource + .getRepository(Stream) + .createQueryBuilder() + .update(Stream) + .set({ + totalAmount: () => `"totalAmount" + :fundAmount`, + }) + .where("id = :id", { id: input.streamId }) + .setParameter("fundAmount", input.amount) + .execute(); + } + + async recordWithdrawal(input: RecordWithdrawalInput): Promise { + await this.dataSource.transaction(async (manager) => { + const result = await manager + .createQueryBuilder() + .insert() + .into(WithdrawalAction) + .values({ + streamId: input.streamId, + recipient: input.recipient, + amount: input.amount, + txHash: input.txHash, + timestamp: input.timestamp, + }) + .execute(); // uuid primary key will always insert, but we should make sure we don't double count if same txHash etc. + + // Actually, wait, if event is replayed, the event index will prevent it. So we don't need orIgnore on uuid. + + await manager + .createQueryBuilder() + .update(Stream) + .set({ + amountWithdrawn: () => `"amountWithdrawn" + :withdrawalAmount`, + }) + .where("id = :id", { id: input.streamId }) + .setParameter("withdrawalAmount", input.amount) + .execute(); + }); + } + + async recordCancel(input: RecordCancelInput): Promise { + await this.dataSource.transaction(async (manager) => { + await manager + .createQueryBuilder() + .insert() + .into(CancelAction) + .values({ + streamId: input.streamId, + canceler: input.canceler, + txHash: input.txHash, + timestamp: input.timestamp, + }) + .execute(); + + await manager + .createQueryBuilder() + .update(Stream) + .set({ + canceled: true, + }) + .where("id = :id", { id: input.streamId }) + .execute(); + }); + } +} diff --git a/indexer/streams/src/handlers/index.ts b/indexer/streams/src/handlers/index.ts index edd56d9..f5875e6 100644 --- a/indexer/streams/src/handlers/index.ts +++ b/indexer/streams/src/handlers/index.ts @@ -1,9 +1,6 @@ -export { streamFundedHandler } from "./stream-funded.handler.js"; -export { streamWithdrawalHandler } from "./stream-withdrawal.handler.js"; -export { streamCancelHandler } from "./stream-cancel.handler.js"; -export { - handleStreamCreated, - parseStreamCreatedPayload, - STREAM_CREATED_TOPIC, -} from "./streamCreated.js"; +export { createStreamFundedHandler } from "./stream-funded.handler.js"; +export { createStreamWithdrawalHandler } from "./stream-withdrawal.handler.js"; +export { createStreamCancelHandler } from "./stream-cancel.handler.js"; +export { createStreamCreatedHandler } from "./stream-created.handler.js"; +export { type StreamHandlerDeps, type EventIdentityStore, deriveEventIndex } from "./persistence.js"; export * from "./types.js"; diff --git a/indexer/streams/src/handlers/persistence.ts b/indexer/streams/src/handlers/persistence.ts new file mode 100644 index 0000000..95328e2 --- /dev/null +++ b/indexer/streams/src/handlers/persistence.ts @@ -0,0 +1,29 @@ +import type { SorobanEventInput } from "@fundable-indexer/common"; +import type { StreamPersistence } from "../db/repository.js"; + +export interface EventIdentityStore { + isEventProcessed( + contractId: string, + ledgerNumber: number, + txHash: string, + eventIndex: number, + ): Promise; + recordEventProcessed( + contractId: string, + ledgerNumber: number, + txHash: string, + eventIndex: number, + ): Promise; +} + +export interface StreamHandlerDeps { + streams: StreamPersistence; + events: EventIdentityStore; +} + +export function deriveEventIndex(event: SorobanEventInput): number { + const segments = event.id.split("-"); + const last = segments[segments.length - 1]; + const parsed = Number.parseInt(last ?? "", 10); + return Number.isFinite(parsed) ? parsed : 0; +} diff --git a/indexer/streams/src/handlers/stream-cancel.handler.ts b/indexer/streams/src/handlers/stream-cancel.handler.ts index 6c01c2b..9f2e1e9 100644 --- a/indexer/streams/src/handlers/stream-cancel.handler.ts +++ b/indexer/streams/src/handlers/stream-cancel.handler.ts @@ -1,35 +1,53 @@ import type { EventHandler, HandlerResult, SorobanEventInput } from "@fundable-indexer/common"; +import { type StreamHandlerDeps, deriveEventIndex } from "./persistence.js"; import { parseStreamCancel } from "./types.js"; -export const streamCancelHandler: EventHandler = async ( - event: SorobanEventInput, -): Promise => { - try { - const payload = parseStreamCancel(event.data); +export const createStreamCancelHandler = (deps: StreamHandlerDeps): EventHandler => { + return async (event: SorobanEventInput): Promise => { + try { + const payload = parseStreamCancel(event.data); - if (!payload.streamId) { - return { ok: false, error: "Missing streamId in cancel event", retriable: false }; - } + if (!payload.streamId) return { ok: false, error: "Missing streamId in cancel event", retriable: false }; + if (!payload.cancelledBy) return { ok: false, error: "Missing cancelledBy in cancel event", retriable: false }; + if (!payload.transactionHash) return { ok: false, error: "Missing transactionHash in cancel event", retriable: false }; - if (!payload.cancelledBy) { - return { ok: false, error: "Missing cancelledBy in cancel event", retriable: false }; - } + const eventIndex = deriveEventIndex(event); + const alreadyProcessed = await deps.events.isEventProcessed( + event.contractId, + event.ledger, + payload.transactionHash, + eventIndex, + ); + if (alreadyProcessed) return { ok: true }; - if (!payload.transactionHash) { - return { ok: false, error: "Missing transactionHash in cancel event", retriable: false }; - } + // Convert ISO string to unix timestamp in seconds for the bigint column + const timestamp = String(Math.floor(Date.parse(event.ledgerClosedAt) / 1000) || 0); + + await deps.streams.recordCancel({ + streamId: payload.streamId, + canceler: payload.cancelledBy, + txHash: payload.transactionHash, + timestamp, + }); - // TODO(#32): update stream status to CANCELLED via repository once DB schema is merged - console.info( - `[stream-cancel] streamId=${payload.streamId} cancelledBy=${payload.cancelledBy} senderBalance=${payload.senderBalance} ledger=${event.ledger}`, - ); - - return { ok: true }; - } catch (err) { - return { - ok: false, - error: err instanceof Error ? err.message : String(err), - retriable: true, - }; - } + await deps.events.recordEventProcessed( + event.contractId, + event.ledger, + payload.transactionHash, + eventIndex, + ); + + console.info( + `[stream-cancel] streamId=${payload.streamId} cancelledBy=${payload.cancelledBy} ledger=${event.ledger}`, + ); + + return { ok: true }; + } catch (err) { + return { + ok: false, + error: err instanceof Error ? err.message : String(err), + retriable: true, + }; + } + }; }; diff --git a/indexer/streams/src/handlers/stream-created.handler.ts b/indexer/streams/src/handlers/stream-created.handler.ts new file mode 100644 index 0000000..aa206f3 --- /dev/null +++ b/indexer/streams/src/handlers/stream-created.handler.ts @@ -0,0 +1,58 @@ +import type { EventHandler, HandlerResult, SorobanEventInput } from "@fundable-indexer/common"; +import { type StreamHandlerDeps, deriveEventIndex } from "./persistence.js"; +import { parseStreamCreated } from "./types.js"; + +export const createStreamCreatedHandler = (deps: StreamHandlerDeps): EventHandler => { + return async (event: SorobanEventInput): Promise => { + try { + const payload = parseStreamCreated(event.data); + + if (!payload.streamId) return { ok: false, error: "Missing streamId in created event", retriable: false }; + if (!payload.sender) return { ok: false, error: "Missing sender in created event", retriable: false }; + if (!payload.recipient) return { ok: false, error: "Missing recipient in created event", retriable: false }; + if (!payload.amount) return { ok: false, error: "Missing amount in created event", retriable: false }; + if (!payload.startTime) return { ok: false, error: "Missing startTime in created event", retriable: false }; + if (!payload.endTime) return { ok: false, error: "Missing endTime in created event", retriable: false }; + if (!payload.transactionHash) return { ok: false, error: "Missing transactionHash in created event", retriable: false }; + + const eventIndex = deriveEventIndex(event); + + const alreadyProcessed = await deps.events.isEventProcessed( + event.contractId, + event.ledger, + payload.transactionHash, + eventIndex, + ); + if (alreadyProcessed) return { ok: true }; + + await deps.streams.createStream({ + id: payload.streamId, + sender: payload.sender, + recipient: payload.recipient, + token: payload.token || "", // Provide empty string if missing, though it's likely present + totalAmount: payload.amount, + startTime: payload.startTime, + endTime: payload.endTime, + }); + + await deps.events.recordEventProcessed( + event.contractId, + event.ledger, + payload.transactionHash, + eventIndex, + ); + + console.info( + `[stream-created] streamId=${payload.streamId} amount=${payload.amount} ledger=${event.ledger}`, + ); + + return { ok: true }; + } catch (err) { + return { + ok: false, + error: err instanceof Error ? err.message : String(err), + retriable: true, + }; + } + }; +}; diff --git a/indexer/streams/src/handlers/stream-funded.handler.ts b/indexer/streams/src/handlers/stream-funded.handler.ts index 893663e..afa4ad9 100644 --- a/indexer/streams/src/handlers/stream-funded.handler.ts +++ b/indexer/streams/src/handlers/stream-funded.handler.ts @@ -1,43 +1,50 @@ import type { EventHandler, HandlerResult, SorobanEventInput } from "@fundable-indexer/common"; +import { type StreamHandlerDeps, deriveEventIndex } from "./persistence.js"; import { parseStreamFunded } from "./types.js"; -export const streamFundedHandler: EventHandler = async ( - event: SorobanEventInput, -): Promise => { - try { - const payload = parseStreamFunded(event.data); - - if (!payload.streamId) { - return { ok: false, error: "Missing streamId in funded event", retriable: false }; - } - - if (!payload.sender) { - return { ok: false, error: "Missing sender in funded event", retriable: false }; +export const createStreamFundedHandler = (deps: StreamHandlerDeps): EventHandler => { + return async (event: SorobanEventInput): Promise => { + try { + const payload = parseStreamFunded(event.data); + + if (!payload.streamId) return { ok: false, error: "Missing streamId in funded event", retriable: false }; + if (!payload.sender) return { ok: false, error: "Missing sender in funded event", retriable: false }; + if (!payload.token) return { ok: false, error: "Missing token in funded event", retriable: false }; + if (!payload.amount) return { ok: false, error: "Missing amount in funded event", retriable: false }; + if (!payload.transactionHash) return { ok: false, error: "Missing transactionHash in funded event", retriable: false }; + + const eventIndex = deriveEventIndex(event); + const alreadyProcessed = await deps.events.isEventProcessed( + event.contractId, + event.ledger, + payload.transactionHash, + eventIndex, + ); + if (alreadyProcessed) return { ok: true }; + + await deps.streams.fundStream({ + streamId: payload.streamId, + amount: payload.amount, + }); + + await deps.events.recordEventProcessed( + event.contractId, + event.ledger, + payload.transactionHash, + eventIndex, + ); + + console.info( + `[stream-funded] streamId=${payload.streamId} amount=${payload.amount} token=${payload.token} ledger=${event.ledger}`, + ); + + return { ok: true }; + } catch (err) { + return { + ok: false, + error: err instanceof Error ? err.message : String(err), + retriable: true, + }; } - - if (!payload.token) { - return { ok: false, error: "Missing token in funded event", retriable: false }; - } - - if (!payload.amount) { - return { ok: false, error: "Missing amount in funded event", retriable: false }; - } - - if (!payload.transactionHash) { - return { ok: false, error: "Missing transactionHash in funded event", retriable: false }; - } - - // TODO(#32): persist deposit via stream repository once DB schema is merged - console.info( - `[stream-funded] streamId=${payload.streamId} amount=${payload.amount} token=${payload.token} ledger=${event.ledger}`, - ); - - return { ok: true }; - } catch (err) { - return { - ok: false, - error: err instanceof Error ? err.message : String(err), - retriable: true, - }; - } + }; }; diff --git a/indexer/streams/src/handlers/stream-handlers.test.ts b/indexer/streams/src/handlers/stream-handlers.test.ts index 245831b..8d9d27c 100644 --- a/indexer/streams/src/handlers/stream-handlers.test.ts +++ b/indexer/streams/src/handlers/stream-handlers.test.ts @@ -1,9 +1,12 @@ -import { describe, expect, test } from "vitest"; +import { beforeEach, describe, expect, test, vi } from "vitest"; import type { SorobanEventInput } from "@fundable-indexer/common"; -import { streamCancelHandler } from "./stream-cancel.handler.js"; -import { streamFundedHandler } from "./stream-funded.handler.js"; -import { streamWithdrawalHandler } from "./stream-withdrawal.handler.js"; +import { createStreamCancelHandler } from "./stream-cancel.handler.js"; +import { createStreamFundedHandler } from "./stream-funded.handler.js"; +import { createStreamWithdrawalHandler } from "./stream-withdrawal.handler.js"; +import { createStreamCreatedHandler } from "./stream-created.handler.js"; +import type { StreamHandlerDeps, EventIdentityStore } from "./persistence.js"; +import type { StreamPersistence } from "../db/repository.js"; const baseEvent: SorobanEventInput = { contractId: "CSTREAM123", @@ -15,98 +18,215 @@ const baseEvent: SorobanEventInput = { pagingToken: "paging-2", }; -describe("streamFundedHandler", () => { - test("returns ok for valid funded payload", async () => { - const event: SorobanEventInput = { - ...baseEvent, - topic: ["stream_funded"], - data: { - stream_id: "stream-1", - sender: "GSENDER", - amount: "5000", - token: "USDC", - tx_hash: "abc123", - }, +describe("Stream Handlers", () => { + let mockStreams: import("vitest").Mocked; + let mockEvents: import("vitest").Mocked; + let deps: StreamHandlerDeps; + + beforeEach(() => { + mockStreams = { + createStream: vi.fn(), + fundStream: vi.fn(), + recordWithdrawal: vi.fn(), + recordCancel: vi.fn(), + } as any; + + mockEvents = { + isEventProcessed: vi.fn().mockResolvedValue(false), + recordEventProcessed: vi.fn().mockResolvedValue(true), + } as any; + + deps = { + streams: mockStreams, + events: mockEvents, }; - - const result = await streamFundedHandler(event); - expect(result).toEqual({ ok: true }); }); - test("returns error when streamId is missing", async () => { - const event: SorobanEventInput = { - ...baseEvent, - data: { amount: "100", token: "XLM", sender: "G123" }, - }; - - const result = await streamFundedHandler(event); - expect(result).toMatchObject({ ok: false, retriable: false }); + describe("streamCreatedHandler", () => { + test("returns ok for valid created payload", async () => { + const handler = createStreamCreatedHandler(deps); + const event: SorobanEventInput = { + ...baseEvent, + topic: ["stream_created"], + data: { + stream_id: "stream-1", + sender: "GSENDER", + recipient: "GRECIPIENT", + amount: "1000", + start_time: "10000", + end_time: "20000", + tx_hash: "tx123", + }, + }; + + const result = await handler(event); + expect(result).toEqual({ ok: true }); + expect(mockStreams.createStream).toHaveBeenCalledWith({ + id: "stream-1", + sender: "GSENDER", + recipient: "GRECIPIENT", + token: "N/A", + totalAmount: "1000", + startTime: "10000", + endTime: "20000", + }); + expect(mockEvents.recordEventProcessed).toHaveBeenCalled(); + }); + + test("skips processing if event is already processed", async () => { + mockEvents.isEventProcessed.mockResolvedValueOnce(true); + const handler = createStreamCreatedHandler(deps); + const event: SorobanEventInput = { + ...baseEvent, + data: { + stream_id: "stream-1", + sender: "GSENDER", + recipient: "GRECIPIENT", + amount: "1000", + start_time: "10000", + end_time: "20000", + tx_hash: "tx123", + }, + }; + + const result = await handler(event); + expect(result).toEqual({ ok: true }); + expect(mockStreams.createStream).not.toHaveBeenCalled(); + }); }); - test("handles unexpected data shape without throwing", async () => { - const event: SorobanEventInput = { - ...baseEvent, - data: null, - }; - - const result = await streamFundedHandler(event); - expect(result.ok).toBe(false); + describe("streamFundedHandler", () => { + test("returns ok for valid funded payload", async () => { + const handler = createStreamFundedHandler(deps); + const event: SorobanEventInput = { + ...baseEvent, + topic: ["stream_funded"], + data: { + stream_id: "stream-1", + sender: "GSENDER", + amount: "5000", + token: "USDC", + tx_hash: "abc123", + }, + }; + + const result = await handler(event); + expect(result).toEqual({ ok: true }); + expect(mockStreams.fundStream).toHaveBeenCalledWith({ + streamId: "stream-1", + amount: "5000", + }); + expect(mockEvents.recordEventProcessed).toHaveBeenCalled(); + }); + + test("skips processing if event is already processed", async () => { + mockEvents.isEventProcessed.mockResolvedValueOnce(true); + const handler = createStreamFundedHandler(deps); + const event: SorobanEventInput = { + ...baseEvent, + data: { + stream_id: "stream-1", + sender: "GSENDER", + amount: "5000", + token: "USDC", + tx_hash: "abc123", + }, + }; + + const result = await handler(event); + expect(result).toEqual({ ok: true }); + expect(mockStreams.fundStream).not.toHaveBeenCalled(); + }); }); -}); -describe("streamWithdrawalHandler", () => { - test("returns ok for valid withdrawal payload", async () => { - const event: SorobanEventInput = { - ...baseEvent, - topic: ["stream_withdrawal"], - data: { - stream_id: "stream-1", + describe("streamWithdrawalHandler", () => { + test("returns ok for valid withdrawal payload", async () => { + const handler = createStreamWithdrawalHandler(deps); + const event: SorobanEventInput = { + ...baseEvent, + topic: ["stream_withdrawal"], + data: { + stream_id: "stream-1", + recipient: "GRECIPIENT", + amount: "250", + tx_hash: "def456", + }, + }; + + const result = await handler(event); + expect(result).toEqual({ ok: true }); + expect(mockStreams.recordWithdrawal).toHaveBeenCalledWith({ + streamId: "stream-1", recipient: "GRECIPIENT", amount: "250", - tx_hash: "def456", - }, - }; - - const result = await streamWithdrawalHandler(event); - expect(result).toEqual({ ok: true }); - }); - - test("returns error when streamId is missing", async () => { - const event: SorobanEventInput = { - ...baseEvent, - data: { recipient: "G123", amount: "50" }, - }; - - const result = await streamWithdrawalHandler(event); - expect(result).toMatchObject({ ok: false, retriable: false }); - }); -}); - -describe("streamCancelHandler", () => { - test("returns ok for valid cancel payload", async () => { - const event: SorobanEventInput = { - ...baseEvent, - topic: ["stream_cancel"], - data: { - stream_id: "stream-1", - cancelled_by: "GSENDER", - sender_balance: "4750", - recipient_balance: "250", - tx_hash: "ghi789", - }, - }; - - const result = await streamCancelHandler(event); - expect(result).toEqual({ ok: true }); + txHash: "def456", + timestamp: "1717200000", + }); + expect(mockEvents.recordEventProcessed).toHaveBeenCalled(); + }); + + test("skips processing if event is already processed", async () => { + mockEvents.isEventProcessed.mockResolvedValueOnce(true); + const handler = createStreamWithdrawalHandler(deps); + const event: SorobanEventInput = { + ...baseEvent, + data: { + stream_id: "stream-1", + recipient: "GRECIPIENT", + amount: "250", + tx_hash: "def456", + }, + }; + + const result = await handler(event); + expect(result).toEqual({ ok: true }); + expect(mockStreams.recordWithdrawal).not.toHaveBeenCalled(); + }); }); - test("returns error when streamId is missing", async () => { - const event: SorobanEventInput = { - ...baseEvent, - data: { cancelled_by: "G123" }, - }; - - const result = await streamCancelHandler(event); - expect(result).toMatchObject({ ok: false, retriable: false }); + describe("streamCancelHandler", () => { + test("returns ok for valid cancel payload", async () => { + const handler = createStreamCancelHandler(deps); + const event: SorobanEventInput = { + ...baseEvent, + topic: ["stream_cancel"], + data: { + stream_id: "stream-1", + cancelled_by: "GSENDER", + sender_balance: "4750", + recipient_balance: "250", + tx_hash: "ghi789", + }, + }; + + const result = await handler(event); + expect(result).toEqual({ ok: true }); + expect(mockStreams.recordCancel).toHaveBeenCalledWith({ + streamId: "stream-1", + canceler: "GSENDER", + txHash: "ghi789", + timestamp: "1717200000", + }); + expect(mockEvents.recordEventProcessed).toHaveBeenCalled(); + }); + + test("skips processing if event is already processed", async () => { + mockEvents.isEventProcessed.mockResolvedValueOnce(true); + const handler = createStreamCancelHandler(deps); + const event: SorobanEventInput = { + ...baseEvent, + data: { + stream_id: "stream-1", + cancelled_by: "GSENDER", + sender_balance: "4750", + recipient_balance: "250", + tx_hash: "ghi789", + }, + }; + + const result = await handler(event); + expect(result).toEqual({ ok: true }); + expect(mockStreams.recordCancel).not.toHaveBeenCalled(); + }); }); }); diff --git a/indexer/streams/src/handlers/stream-withdrawal.handler.ts b/indexer/streams/src/handlers/stream-withdrawal.handler.ts index d258ff4..d1817c5 100644 --- a/indexer/streams/src/handlers/stream-withdrawal.handler.ts +++ b/indexer/streams/src/handlers/stream-withdrawal.handler.ts @@ -1,39 +1,52 @@ import type { EventHandler, HandlerResult, SorobanEventInput } from "@fundable-indexer/common"; +import { type StreamHandlerDeps, deriveEventIndex } from "./persistence.js"; import { parseStreamWithdrawal } from "./types.js"; -export const streamWithdrawalHandler: EventHandler = async ( - event: SorobanEventInput, -): Promise => { - try { - const payload = parseStreamWithdrawal(event.data); +export const createStreamWithdrawalHandler = (deps: StreamHandlerDeps): EventHandler => { + return async (event: SorobanEventInput): Promise => { + try { + const payload = parseStreamWithdrawal(event.data); - if (!payload.streamId) { - return { ok: false, error: "Missing streamId in withdrawal event", retriable: false }; - } + if (!payload.streamId) return { ok: false, error: "Missing streamId in withdrawal event", retriable: false }; + if (!payload.recipient) return { ok: false, error: "Missing recipient in withdrawal event", retriable: false }; + if (!payload.amount) return { ok: false, error: "Missing amount in withdrawal event", retriable: false }; + if (!payload.transactionHash) return { ok: false, error: "Missing transactionHash in withdrawal event", retriable: false }; - if (!payload.recipient) { - return { ok: false, error: "Missing recipient in withdrawal event", retriable: false }; - } + const eventIndex = deriveEventIndex(event); + const alreadyProcessed = await deps.events.isEventProcessed( + event.contractId, + event.ledger, + payload.transactionHash, + eventIndex, + ); + if (alreadyProcessed) return { ok: true }; - if (!payload.amount) { - return { ok: false, error: "Missing amount in withdrawal event", retriable: false }; - } + await deps.streams.recordWithdrawal({ + streamId: payload.streamId, + recipient: payload.recipient, + amount: payload.amount, + txHash: payload.transactionHash, + timestamp: String(Math.floor(Date.parse(event.ledgerClosedAt) / 1000) || 0), + }); - if (!payload.transactionHash) { - return { ok: false, error: "Missing transactionHash in withdrawal event", retriable: false }; - } + await deps.events.recordEventProcessed( + event.contractId, + event.ledger, + payload.transactionHash, + eventIndex, + ); - // TODO(#32): record withdrawal action via repository once DB schema is merged - console.info( - `[stream-withdrawal] streamId=${payload.streamId} recipient=${payload.recipient} amount=${payload.amount} ledger=${event.ledger}`, - ); + console.info( + `[stream-withdrawal] streamId=${payload.streamId} recipient=${payload.recipient} amount=${payload.amount} ledger=${event.ledger}`, + ); - return { ok: true }; - } catch (err) { - return { - ok: false, - error: err instanceof Error ? err.message : String(err), - retriable: true, - }; - } + return { ok: true }; + } catch (err) { + return { + ok: false, + error: err instanceof Error ? err.message : String(err), + retriable: true, + }; + } + }; }; diff --git a/indexer/streams/src/handlers/streamCreated.test.ts b/indexer/streams/src/handlers/streamCreated.test.ts deleted file mode 100644 index e36be52..0000000 --- a/indexer/streams/src/handlers/streamCreated.test.ts +++ /dev/null @@ -1,179 +0,0 @@ -import { describe, expect, test } from "vitest"; - -import { - STREAM_CREATED_TOPIC, - getEventIdentity, - handleStreamCreated, - mapStreamCreatedToRecord, - parseStreamCreatedPayload, -} from "./streamCreated.js"; -import type { StreamCreatedEvent } from "./types.js"; - -function createMockEvent(overrides?: Partial): StreamCreatedEvent { - return { - contractId: "0x123", - ledger: 12345, - txHash: "0xabc", - eventIndex: 0, - topics: [STREAM_CREATED_TOPIC], - data: JSON.stringify({ - streamId: "stream-1", - sender: "0xsender", - recipient: "0xrecipient", - amount: "1000000000", - startTime: "1000000", - endTime: "2000000", - }), - ...overrides, - }; -} - -const mockPayload = { - streamId: "stream-1", - sender: "0xsender", - recipient: "0xrecipient", - amount: "1000000000", - startTime: "1000000", - endTime: "2000000", -}; - -describe("parseStreamCreatedPayload", () => { - test("parses a valid StreamCreated event", () => { - const event = createMockEvent(); - const payload = parseStreamCreatedPayload(event); - - expect(payload).toEqual(mockPayload); - }); - - test("throws when event topic does not match", () => { - const event = createMockEvent({ topics: ["WrongTopic"] }); - - expect(() => parseStreamCreatedPayload(event)).toThrow( - "Expected StreamCreated event topic, got WrongTopic", - ); - }); - - test("throws when topics array is empty", () => { - const event = createMockEvent({ topics: [] }); - - expect(() => parseStreamCreatedPayload(event)).toThrow( - "Expected StreamCreated event topic, got undefined", - ); - }); - - test("throws on invalid JSON data", () => { - const event = createMockEvent({ data: "not-json" }); - - expect(() => parseStreamCreatedPayload(event)).toThrow( - "Failed to parse event data: invalid JSON", - ); - }); - - test("throws on missing streamId field", () => { - const { streamId: _, ...partial } = mockPayload; - const event = createMockEvent({ data: JSON.stringify(partial) }); - - expect(() => parseStreamCreatedPayload(event)).toThrow( - 'Invalid payload: "streamId" must be a non-empty string', - ); - }); - - test("throws on non-string amount field", () => { - const event = createMockEvent({ - data: JSON.stringify({ ...mockPayload, amount: 12345 }), - }); - - expect(() => parseStreamCreatedPayload(event)).toThrow( - 'Invalid payload: "amount" must be a non-empty string', - ); - }); - - test("throws on empty string recipient", () => { - const event = createMockEvent({ - data: JSON.stringify({ ...mockPayload, recipient: "" }), - }); - - expect(() => parseStreamCreatedPayload(event)).toThrow( - 'Invalid payload: "recipient" must be a non-empty string', - ); - }); -}); - -describe("getEventIdentity", () => { - test("produces a deterministic identity string", () => { - const event = createMockEvent(); - const identity = getEventIdentity(event); - - expect(identity).toBe("0x123:12345:0xabc:0"); - }); - - test("changes when any identity field changes", () => { - const base = createMockEvent(); - const differentContract = createMockEvent({ contractId: "0x456" }); - const differentLedger = createMockEvent({ ledger: 99999 }); - const differentTxHash = createMockEvent({ txHash: "0xdef" }); - const differentIndex = createMockEvent({ eventIndex: 1 }); - - const baseId = getEventIdentity(base); - expect(getEventIdentity(differentContract)).not.toBe(baseId); - expect(getEventIdentity(differentLedger)).not.toBe(baseId); - expect(getEventIdentity(differentTxHash)).not.toBe(baseId); - expect(getEventIdentity(differentIndex)).not.toBe(baseId); - }); -}); - -describe("mapStreamCreatedToRecord", () => { - test("maps payload and event to a Stream record", () => { - const event = createMockEvent(); - const record = mapStreamCreatedToRecord(mockPayload, event); - - expect(record).toEqual({ - id: "stream-1", - sender: "0xsender", - recipient: "0xrecipient", - amount: "1000000000", - startTime: "1000000", - endTime: "2000000", - contractId: "0x123", - ledger: 12345, - txHash: "0xabc", - eventIndex: 0, - }); - }); -}); - -describe("handleStreamCreated", () => { - test("returns stream record and identity", () => { - const event = createMockEvent(); - const result = handleStreamCreated(event); - - expect(result.stream.id).toBe("stream-1"); - expect(result.stream.contractId).toBe("0x123"); - expect(result.stream.ledger).toBe(12345); - expect(result.identity).toBe("0x123:12345:0xabc:0"); - }); - - test("is idempotent (same input produces same output)", () => { - const event = createMockEvent(); - const result1 = handleStreamCreated(event); - const result2 = handleStreamCreated(event); - - expect(result1).toEqual(result2); - }); - - test("handles different stream IDs correctly", () => { - const event1 = createMockEvent({ - data: JSON.stringify({ ...mockPayload, streamId: "stream-1" }), - }); - const event2 = createMockEvent({ - data: JSON.stringify({ ...mockPayload, streamId: "stream-2" }), - }); - - const result1 = handleStreamCreated(event1); - const result2 = handleStreamCreated(event2); - - expect(result1.stream.id).toBe("stream-1"); - expect(result2.stream.id).toBe("stream-2"); - expect(result1.stream.id).not.toBe(result2.stream.id); - }); -}); diff --git a/indexer/streams/src/handlers/streamCreated.ts b/indexer/streams/src/handlers/streamCreated.ts deleted file mode 100644 index 33df1a0..0000000 --- a/indexer/streams/src/handlers/streamCreated.ts +++ /dev/null @@ -1,88 +0,0 @@ -import type { StreamCreatedEvent, StreamRecord } from "./types.js"; - -export const STREAM_CREATED_TOPIC = "StreamCreated"; - -export interface StreamCreatedPayload { - streamId: string; - sender: string; - recipient: string; - amount: string; - startTime: string; - endTime: string; -} - -export function parseStreamCreatedPayload(event: StreamCreatedEvent): StreamCreatedPayload { - const eventName = event.topics[0]; - if (!eventName || eventName !== STREAM_CREATED_TOPIC) { - throw new Error( - `Expected ${STREAM_CREATED_TOPIC} event topic, got ${eventName ?? "undefined"}`, - ); - } - - let parsed: Record; - try { - parsed = JSON.parse(event.data) as Record; - } catch { - throw new Error("Failed to parse event data: invalid JSON"); - } - - const requiredFields = [ - "streamId", - "sender", - "recipient", - "amount", - "startTime", - "endTime", - ] as const; - - for (const field of requiredFields) { - const value = parsed[field]; - if (typeof value !== "string" || value.length === 0) { - throw new Error( - `Invalid payload: "${field}" must be a non-empty string, got ${typeof value === "string" ? "empty string" : typeof value}`, - ); - } - } - - return { - streamId: parsed.streamId as string, - sender: parsed.sender as string, - recipient: parsed.recipient as string, - amount: parsed.amount as string, - startTime: parsed.startTime as string, - endTime: parsed.endTime as string, - }; -} - -export function getEventIdentity(event: StreamCreatedEvent): string { - return `${event.contractId}:${event.ledger}:${event.txHash}:${event.eventIndex}`; -} - -export function mapStreamCreatedToRecord( - payload: StreamCreatedPayload, - event: StreamCreatedEvent, -): StreamRecord { - return { - id: payload.streamId, - sender: payload.sender, - recipient: payload.recipient, - amount: payload.amount, - startTime: payload.startTime, - endTime: payload.endTime, - contractId: event.contractId, - ledger: event.ledger, - txHash: event.txHash, - eventIndex: event.eventIndex, - }; -} - -export function handleStreamCreated(event: StreamCreatedEvent): { - stream: StreamRecord; - identity: string; -} { - const payload = parseStreamCreatedPayload(event); - const stream = mapStreamCreatedToRecord(payload, event); - const identity = getEventIdentity(event); - - return { stream, identity }; -} diff --git a/indexer/streams/src/handlers/types.ts b/indexer/streams/src/handlers/types.ts index ad59caf..5154d06 100644 --- a/indexer/streams/src/handlers/types.ts +++ b/indexer/streams/src/handlers/types.ts @@ -84,3 +84,28 @@ export function parseStreamCancel(data: unknown): StreamCancelPayload { transactionHash: str(d.transactionHash ?? d.tx_hash), }; } + +export interface StreamCreatedPayloadParsed { + streamId: string | undefined; + sender: string | undefined; + recipient: string | undefined; + token: string | undefined; + amount: string | undefined; + startTime: string | undefined; + endTime: string | undefined; + transactionHash: string | undefined; +} + +export function parseStreamCreated(data: unknown): StreamCreatedPayloadParsed { + const d = record(data); + return { + streamId: str(d.streamId ?? d.stream_id), + sender: str(d.sender), + recipient: str(d.recipient), + token: str(d.token), + amount: str(d.amount), + startTime: str(d.startTime ?? d.start_time), + endTime: str(d.endTime ?? d.end_time), + transactionHash: str(d.transactionHash ?? d.tx_hash), + }; +} diff --git a/indexer/streams/src/index.test.ts b/indexer/streams/src/index.test.ts deleted file mode 100644 index cab3906..0000000 --- a/indexer/streams/src/index.test.ts +++ /dev/null @@ -1,13 +0,0 @@ -import { describe, expect, test } from "vitest"; - -import { streamsPackage } from "./index.js"; - -describe("streams package", () => { - test("exposes package identity", () => { - expect(streamsPackage).toEqual({ - name: "@fundable-indexer/streams", - role: "payment-stream-indexer", - common: "@fundable-indexer/common", - }); - }); -}); diff --git a/indexer/streams/src/index.ts b/indexer/streams/src/index.ts index 731f74f..c7e58d2 100644 --- a/indexer/streams/src/index.ts +++ b/indexer/streams/src/index.ts @@ -9,9 +9,5 @@ export const streamsPackage = { export { Stream } from "./db/entity/Stream.js"; export { WithdrawalAction } from "./db/entity/WithdrawalAction.js"; export { CancelAction } from "./db/entity/CancelAction.js"; +export { StreamRepository, type StreamPersistence } from "./db/repository.js"; export * from "./handlers/index.js"; -export { - handleStreamCreated, - parseStreamCreatedPayload, - STREAM_CREATED_TOPIC, -} from "./handlers/streamCreated.js";