diff --git a/indexer/streams/src/db/repository.ts b/indexer/streams/src/db/repository.ts new file mode 100644 index 0000000..f6f328b --- /dev/null +++ b/indexer/streams/src/db/repository.ts @@ -0,0 +1,227 @@ +import type { DataSource, Repository } from "typeorm"; +import { CancelAction } from "./entity/CancelAction.js"; +import { Stream } from "./entity/Stream.js"; +import { WithdrawalAction } from "./entity/WithdrawalAction.js"; + +export interface StreamCreateInput { + id: string; + sender: string; + recipient: string; + token: string; + totalAmount: string; + startTime: string; + endTime: string; +} + +export interface StreamWriteService { + createStream(input: StreamCreateInput, eventKey?: string): Promise; + fundStream(streamId: string, amount: string, eventKey?: string): Promise; + recordWithdrawal( + streamId: string, + recipient: string, + amount: string, + txHash: string, + timestamp: string, + eventKey?: string, + ): Promise; + recordCancel( + streamId: string, + canceler: string, + txHash: string, + timestamp: string, + eventKey?: string, + ): Promise; +} + +export class StreamWriteServiceImpl implements StreamWriteService { + private processedEventKeys = new Set(); + + constructor(private readonly repository: StreamRepository) {} + + async createStream(input: StreamCreateInput, eventKey?: string): Promise { + if (eventKey && this.processedEventKeys.has(eventKey)) { + return; + } + + await this.repository.createStream(input); + + if (eventKey) { + this.processedEventKeys.add(eventKey); + } + } + + async fundStream(streamId: string, amount: string, eventKey?: string): Promise { + if (eventKey && this.processedEventKeys.has(eventKey)) { + return; + } + + await this.repository.fundStream(streamId, amount); + + if (eventKey) { + this.processedEventKeys.add(eventKey); + } + } + + async recordWithdrawal( + streamId: string, + recipient: string, + amount: string, + txHash: string, + timestamp: string, + eventKey?: string, + ): Promise { + if (eventKey && this.processedEventKeys.has(eventKey)) { + return; + } + + await this.repository.recordWithdrawal(streamId, recipient, amount, txHash, timestamp); + + if (eventKey) { + this.processedEventKeys.add(eventKey); + } + } + + async recordCancel( + streamId: string, + canceler: string, + txHash: string, + timestamp: string, + eventKey?: string, + ): Promise { + if (eventKey && this.processedEventKeys.has(eventKey)) { + return; + } + + await this.repository.recordCancel(streamId, canceler, txHash, timestamp); + + if (eventKey) { + this.processedEventKeys.add(eventKey); + } + } +} + +export function createStreamWriteService(repository: StreamRepository): StreamWriteService { + return new StreamWriteServiceImpl(repository); +} + +export class StreamRepository { + private streamRepo: Repository; + private withdrawalRepo: Repository; + private cancelRepo: Repository; + + constructor(private dataSource: DataSource) { + this.streamRepo = this.dataSource.getRepository(Stream); + this.withdrawalRepo = this.dataSource.getRepository(WithdrawalAction); + this.cancelRepo = this.dataSource.getRepository(CancelAction); + } + + async createStream(input: StreamCreateInput): Promise { + await this.streamRepo.upsert( + this.streamRepo.create({ + 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, + }), + ["id"], + ); + } + + async getStreamById(streamId: string): Promise { + return this.streamRepo.findOneBy({ id: streamId }); + } + + async fundStream(streamId: string, amount: string): Promise { + const amountBigInt = BigInt(amount); + if (amountBigInt <= 0n) { + throw new Error("Funding amount must be positive"); + } + + const stream = await this.getStreamById(streamId); + if (!stream) { + throw new Error(`Stream not found: ${streamId}`); + } + + const nextTotalAmount = (BigInt(stream.totalAmount) + amountBigInt).toString(); + await this.streamRepo.update({ id: streamId }, { totalAmount: nextTotalAmount }); + } + + async recordWithdrawal( + streamId: string, + recipient: string, + amount: string, + txHash: string, + timestamp: string, + ): Promise { + const existingWithdrawal = await this.withdrawalRepo.findOne({ + where: { streamId, txHash }, + }); + if (existingWithdrawal) { + return; + } + + const stream = await this.getStreamById(streamId); + if (!stream) { + throw new Error(`Stream not found: ${streamId}`); + } + + const amountBigInt = BigInt(amount); + if (amountBigInt <= 0n) { + throw new Error("Withdrawal amount must be positive"); + } + + await this.dataSource.transaction(async (manager) => { + const withdrawalRepo = manager.getRepository(WithdrawalAction); + const streamRepo = manager.getRepository(Stream); + + await withdrawalRepo.insert({ + streamId, + recipient, + amount, + txHash, + timestamp, + }); + + const nextWithdrawn = (BigInt(stream.amountWithdrawn) + amountBigInt).toString(); + await streamRepo.update({ id: streamId }, { amountWithdrawn: nextWithdrawn }); + }); + } + + async recordCancel( + streamId: string, + canceler: string, + txHash: string, + timestamp: string, + ): Promise { + const existingCancel = await this.cancelRepo.findOne({ + where: { streamId, txHash }, + }); + if (existingCancel) { + return; + } + + const stream = await this.getStreamById(streamId); + if (!stream) { + throw new Error(`Stream not found: ${streamId}`); + } + + await this.dataSource.transaction(async (manager) => { + const cancelRepo = manager.getRepository(CancelAction); + const streamRepo = manager.getRepository(Stream); + + await cancelRepo.insert({ + streamId, + canceler, + txHash, + timestamp, + }); + + await streamRepo.update({ id: streamId }, { canceled: true }); + }); + } +} diff --git a/indexer/streams/src/handlers/stream-cancel.handler.ts b/indexer/streams/src/handlers/stream-cancel.handler.ts index 6c01c2b..51b5282 100644 --- a/indexer/streams/src/handlers/stream-cancel.handler.ts +++ b/indexer/streams/src/handlers/stream-cancel.handler.ts @@ -1,35 +1,51 @@ +import type { + EventHandler, + HandlerResult, + SorobanEventInput, +} from "@fundable-indexer/common"; +import type { StreamWriteService } from "../db/repository.js"; import type { EventHandler, HandlerResult, SorobanEventInput } from "@fundable-indexer/common"; import { parseStreamCancel } from "./types.js"; -export const streamCancelHandler: EventHandler = async ( - event: SorobanEventInput, -): Promise => { - try { - const payload = parseStreamCancel(event.data); +export function createStreamCancelHandler(persistence?: StreamWriteService): 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.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 }; + } + + const eventKey = `${event.contractId}:${event.ledger}:${payload.transactionHash}`; + await persistence?.recordCancel( + payload.streamId, + payload.cancelledBy, + payload.transactionHash, + new Date(event.ledger * 1000).toISOString(), + eventKey, + ); + + console.info( + `[stream-cancel] streamId=${payload.streamId} cancelledBy=${payload.cancelledBy} senderBalance=${payload.senderBalance} ledger=${event.ledger}`, + ); - if (!payload.transactionHash) { - return { ok: false, error: "Missing transactionHash in cancel event", retriable: false }; + return { ok: true }; + } catch (err) { + return { + ok: false, + error: err instanceof Error ? err.message : String(err), + retriable: true, + }; } + }; +} - // 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, - }; - } -}; +export const streamCancelHandler = createStreamCancelHandler(); diff --git a/indexer/streams/src/handlers/stream-funded.handler.ts b/indexer/streams/src/handlers/stream-funded.handler.ts index 893663e..961e55d 100644 --- a/indexer/streams/src/handlers/stream-funded.handler.ts +++ b/indexer/streams/src/handlers/stream-funded.handler.ts @@ -1,43 +1,53 @@ +import type { + EventHandler, + HandlerResult, + SorobanEventInput, +} from "@fundable-indexer/common"; +import type { StreamWriteService } from "../db/repository.js"; import type { EventHandler, HandlerResult, SorobanEventInput } from "@fundable-indexer/common"; 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 }; - } - - 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 }; +export function createStreamFundedHandler(persistence?: StreamWriteService): 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 eventKey = `${event.contractId}:${event.ledger}:${payload.transactionHash}`; + await persistence?.fundStream(payload.streamId, payload.amount, eventKey); + + 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, + }; } + }; +} - // 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, - }; - } -}; +export const streamFundedHandler = createStreamFundedHandler(); diff --git a/indexer/streams/src/handlers/stream-handlers.test.ts b/indexer/streams/src/handlers/stream-handlers.test.ts index 245831b..879db17 100644 --- a/indexer/streams/src/handlers/stream-handlers.test.ts +++ b/indexer/streams/src/handlers/stream-handlers.test.ts @@ -1,9 +1,76 @@ import { describe, expect, test } 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 type { StreamWriteService } from "../db/repository.js"; +import { createStreamCancelHandler } from "./stream-cancel.handler.js"; +import { createStreamFundedHandler } from "./stream-funded.handler.js"; +import { createStreamWithdrawalHandler } from "./stream-withdrawal.handler.js"; + +class StubStreamWriteService implements StreamWriteService { + public createdStreams: Array<{ id: string; amount: string }> = []; + public fundedStreams: Array<{ streamId: string; amount: string }> = []; + public withdrawals: Array<{ streamId: string; amount: string }> = []; + public cancels: Array<{ streamId: string }> = []; + private processedEvents = new Set(); + + private shouldProcess(eventKey?: string): boolean { + if (!eventKey) { + return true; + } + + if (this.processedEvents.has(eventKey)) { + return false; + } + + this.processedEvents.add(eventKey); + return true; + } + + async createStream(input: { id: string; totalAmount: string }, eventKey?: string): Promise { + if (!this.shouldProcess(eventKey)) { + return; + } + + this.createdStreams.push({ id: input.id, amount: input.totalAmount }); + } + + async fundStream(streamId: string, amount: string, eventKey?: string): Promise { + if (!this.shouldProcess(eventKey)) { + return; + } + + this.fundedStreams.push({ streamId, amount }); + } + + async recordWithdrawal( + streamId: string, + _recipient: string, + amount: string, + _txHash: string, + _timestamp: string, + eventKey?: string, + ): Promise { + if (!this.shouldProcess(eventKey)) { + return; + } + + this.withdrawals.push({ streamId, amount }); + } + + async recordCancel( + streamId: string, + _canceler: string, + _txHash: string, + _timestamp: string, + eventKey?: string, + ): Promise { + if (!this.shouldProcess(eventKey)) { + return; + } + + this.cancels.push({ streamId }); + } +} const baseEvent: SorobanEventInput = { contractId: "CSTREAM123", @@ -16,7 +83,9 @@ const baseEvent: SorobanEventInput = { }; describe("streamFundedHandler", () => { - test("returns ok for valid funded payload", async () => { + test("persists funded amount through the stream write service", async () => { + const persistence = new StubStreamWriteService(); + const handler = createStreamFundedHandler(persistence); const event: SorobanEventInput = { ...baseEvent, topic: ["stream_funded"], @@ -29,8 +98,33 @@ describe("streamFundedHandler", () => { }, }; - const result = await streamFundedHandler(event); + const result = await handler(event); expect(result).toEqual({ ok: true }); + expect(persistence.fundedStreams).toEqual([{ streamId: "stream-1", amount: "5000" }]); + }); + + test("does not double-count a duplicate funded replay", async () => { + const persistence = new StubStreamWriteService(); + const handler = createStreamFundedHandler(persistence); + const event: SorobanEventInput = { + ...baseEvent, + topic: ["stream_funded"], + data: { + stream_id: "stream-1", + sender: "GSENDER", + amount: "5000", + token: "USDC", + tx_hash: "abc123", + }, + }; + + await handler(event); + await handler(event); + + expect(persistence.fundedStreams).toHaveLength(1); + }); + + test("returns ok for valid funded payload", async () => { }); test("returns error when streamId is missing", async () => { @@ -39,7 +133,7 @@ describe("streamFundedHandler", () => { data: { amount: "100", token: "XLM", sender: "G123" }, }; - const result = await streamFundedHandler(event); + const result = await createStreamFundedHandler()(event); expect(result).toMatchObject({ ok: false, retriable: false }); }); @@ -49,13 +143,15 @@ describe("streamFundedHandler", () => { data: null, }; - const result = await streamFundedHandler(event); + const result = await createStreamFundedHandler()(event); expect(result.ok).toBe(false); }); }); describe("streamWithdrawalHandler", () => { - test("returns ok for valid withdrawal payload", async () => { + test("persists withdrawal actions through the stream write service", async () => { + const persistence = new StubStreamWriteService(); + const handler = createStreamWithdrawalHandler(persistence); const event: SorobanEventInput = { ...baseEvent, topic: ["stream_withdrawal"], @@ -67,8 +163,32 @@ describe("streamWithdrawalHandler", () => { }, }; - const result = await streamWithdrawalHandler(event); + const result = await handler(event); expect(result).toEqual({ ok: true }); + expect(persistence.withdrawals).toEqual([{ streamId: "stream-1", amount: "250" }]); + }); + + test("does not double-record a duplicate withdrawal replay", async () => { + const persistence = new StubStreamWriteService(); + const handler = createStreamWithdrawalHandler(persistence); + const event: SorobanEventInput = { + ...baseEvent, + topic: ["stream_withdrawal"], + data: { + stream_id: "stream-1", + recipient: "GRECIPIENT", + amount: "250", + tx_hash: "def456", + }, + }; + + await handler(event); + await handler(event); + + expect(persistence.withdrawals).toHaveLength(1); + }); + + test("returns ok for valid withdrawal payload", async () => { }); test("returns error when streamId is missing", async () => { @@ -77,13 +197,15 @@ describe("streamWithdrawalHandler", () => { data: { recipient: "G123", amount: "50" }, }; - const result = await streamWithdrawalHandler(event); + const result = await createStreamWithdrawalHandler()(event); expect(result).toMatchObject({ ok: false, retriable: false }); }); }); describe("streamCancelHandler", () => { - test("returns ok for valid cancel payload", async () => { + test("persists cancel actions through the stream write service", async () => { + const persistence = new StubStreamWriteService(); + const handler = createStreamCancelHandler(persistence); const event: SorobanEventInput = { ...baseEvent, topic: ["stream_cancel"], @@ -96,8 +218,33 @@ describe("streamCancelHandler", () => { }, }; - const result = await streamCancelHandler(event); + const result = await handler(event); expect(result).toEqual({ ok: true }); + expect(persistence.cancels).toEqual([{ streamId: "stream-1" }]); + }); + + test("does not double-record a duplicate cancel replay", async () => { + const persistence = new StubStreamWriteService(); + const handler = createStreamCancelHandler(persistence); + const event: SorobanEventInput = { + ...baseEvent, + topic: ["stream_cancel"], + data: { + stream_id: "stream-1", + cancelled_by: "GSENDER", + sender_balance: "4750", + recipient_balance: "250", + tx_hash: "ghi789", + }, + }; + + await handler(event); + await handler(event); + + expect(persistence.cancels).toHaveLength(1); + }); + + test("returns ok for valid cancel payload", async () => { }); test("returns error when streamId is missing", async () => { @@ -106,7 +253,7 @@ describe("streamCancelHandler", () => { data: { cancelled_by: "G123" }, }; - const result = await streamCancelHandler(event); + const result = await createStreamCancelHandler()(event); expect(result).toMatchObject({ ok: false, retriable: false }); }); }); diff --git a/indexer/streams/src/handlers/stream-withdrawal.handler.ts b/indexer/streams/src/handlers/stream-withdrawal.handler.ts index d258ff4..36a29eb 100644 --- a/indexer/streams/src/handlers/stream-withdrawal.handler.ts +++ b/indexer/streams/src/handlers/stream-withdrawal.handler.ts @@ -1,39 +1,56 @@ +import type { + EventHandler, + HandlerResult, + SorobanEventInput, +} from "@fundable-indexer/common"; +import type { StreamWriteService } from "../db/repository.js"; import type { EventHandler, HandlerResult, SorobanEventInput } from "@fundable-indexer/common"; import { parseStreamWithdrawal } from "./types.js"; -export const streamWithdrawalHandler: EventHandler = async ( - event: SorobanEventInput, -): Promise => { - try { - const payload = parseStreamWithdrawal(event.data); +export function createStreamWithdrawalHandler(persistence?: StreamWriteService): 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.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.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 }; + } + + const eventKey = `${event.contractId}:${event.ledger}:${payload.transactionHash}`; + await persistence?.recordWithdrawal( + payload.streamId, + payload.recipient, + payload.amount, + payload.transactionHash, + new Date(event.ledger * 1000).toISOString(), + eventKey, + ); + + console.info( + `[stream-withdrawal] streamId=${payload.streamId} recipient=${payload.recipient} amount=${payload.amount} ledger=${event.ledger}`, + ); - if (!payload.transactionHash) { - return { ok: false, error: "Missing transactionHash in withdrawal event", retriable: false }; + return { ok: true }; + } catch (err) { + return { + ok: false, + error: err instanceof Error ? err.message : String(err), + retriable: true, + }; } + }; +} - // 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}`, - ); - - return { ok: true }; - } catch (err) { - return { - ok: false, - error: err instanceof Error ? err.message : String(err), - retriable: true, - }; - } -}; +export const streamWithdrawalHandler = createStreamWithdrawalHandler(); diff --git a/indexer/streams/src/handlers/streamCreated.test.ts b/indexer/streams/src/handlers/streamCreated.test.ts index e36be52..4dbbe8e 100644 --- a/indexer/streams/src/handlers/streamCreated.test.ts +++ b/indexer/streams/src/handlers/streamCreated.test.ts @@ -1,5 +1,6 @@ import { describe, expect, test } from "vitest"; +import type { StreamWriteService } from "../db/repository.js"; import { STREAM_CREATED_TOPIC, getEventIdentity, @@ -143,9 +144,9 @@ describe("mapStreamCreatedToRecord", () => { }); describe("handleStreamCreated", () => { - test("returns stream record and identity", () => { + test("returns stream record and identity", async () => { const event = createMockEvent(); - const result = handleStreamCreated(event); + const result = await handleStreamCreated(event); expect(result.stream.id).toBe("stream-1"); expect(result.stream.contractId).toBe("0x123"); @@ -153,15 +154,32 @@ describe("handleStreamCreated", () => { expect(result.identity).toBe("0x123:12345:0xabc:0"); }); - test("is idempotent (same input produces same output)", () => { + test("persists a created stream through the write service", async () => { + const persistence: StreamWriteService = { + createStream: async (input) => { + expect(input.id).toBe("stream-1"); + expect(input.totalAmount).toBe("1000000000"); + }, + fundStream: async () => {}, + recordWithdrawal: async () => {}, + recordCancel: async () => {}, + }; + + const event = createMockEvent(); + await expect(handleStreamCreated(event, persistence)).resolves.toMatchObject({ + stream: { id: "stream-1" }, + }); + }); + + test("is idempotent (same input produces same output)", async () => { const event = createMockEvent(); - const result1 = handleStreamCreated(event); - const result2 = handleStreamCreated(event); + const result1 = await handleStreamCreated(event); + const result2 = await handleStreamCreated(event); expect(result1).toEqual(result2); }); - test("handles different stream IDs correctly", () => { + test("handles different stream IDs correctly", async () => { const event1 = createMockEvent({ data: JSON.stringify({ ...mockPayload, streamId: "stream-1" }), }); @@ -169,8 +187,8 @@ describe("handleStreamCreated", () => { data: JSON.stringify({ ...mockPayload, streamId: "stream-2" }), }); - const result1 = handleStreamCreated(event1); - const result2 = handleStreamCreated(event2); + const result1 = await handleStreamCreated(event1); + const result2 = await handleStreamCreated(event2); expect(result1.stream.id).toBe("stream-1"); expect(result2.stream.id).toBe("stream-2"); diff --git a/indexer/streams/src/handlers/streamCreated.ts b/indexer/streams/src/handlers/streamCreated.ts index 33df1a0..0976af7 100644 --- a/indexer/streams/src/handlers/streamCreated.ts +++ b/indexer/streams/src/handlers/streamCreated.ts @@ -1,3 +1,4 @@ +import type { StreamCreateInput, StreamWriteService } from "../db/repository.js"; import type { StreamCreatedEvent, StreamRecord } from "./types.js"; export const STREAM_CREATED_TOPIC = "StreamCreated"; @@ -76,13 +77,29 @@ export function mapStreamCreatedToRecord( }; } -export function handleStreamCreated(event: StreamCreatedEvent): { +export async function handleStreamCreated( + event: StreamCreatedEvent, + persistence?: StreamWriteService, +): Promise<{ stream: StreamRecord; identity: string; -} { +}> { const payload = parseStreamCreatedPayload(event); const stream = mapStreamCreatedToRecord(payload, event); const identity = getEventIdentity(event); + await persistence?.createStream( + { + id: stream.id, + sender: stream.sender, + recipient: stream.recipient, + token: "", + totalAmount: stream.amount, + startTime: stream.startTime, + endTime: stream.endTime, + } satisfies StreamCreateInput, + identity, + ); + return { stream, identity }; }