Skip to content
Open
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
133 changes: 133 additions & 0 deletions indexer/streams/src/db/repository.ts
Original file line number Diff line number Diff line change
@@ -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<void>;
fundStream(input: FundStreamInput): Promise<void>;
recordWithdrawal(input: RecordWithdrawalInput): Promise<void>;
recordCancel(input: RecordCancelInput): Promise<void>;
}

export class StreamRepository implements StreamPersistence {
constructor(private readonly dataSource: DataSource) {}

async createStream(input: CreateStreamInput): Promise<void> {
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<void> {
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<void> {
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<void> {
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();
});
}
}
13 changes: 5 additions & 8 deletions indexer/streams/src/handlers/index.ts
Original file line number Diff line number Diff line change
@@ -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";
29 changes: 29 additions & 0 deletions indexer/streams/src/handlers/persistence.ts
Original file line number Diff line number Diff line change
@@ -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<boolean>;
recordEventProcessed(
contractId: string,
ledgerNumber: number,
txHash: string,
eventIndex: number,
): Promise<boolean>;
}

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;
}
72 changes: 45 additions & 27 deletions indexer/streams/src/handlers/stream-cancel.handler.ts
Original file line number Diff line number Diff line change
@@ -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<HandlerResult> => {
try {
const payload = parseStreamCancel(event.data);
export const createStreamCancelHandler = (deps: StreamHandlerDeps): EventHandler => {
return async (event: SorobanEventInput): Promise<HandlerResult> => {
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,
};
}
};
};
58 changes: 58 additions & 0 deletions indexer/streams/src/handlers/stream-created.handler.ts
Original file line number Diff line number Diff line change
@@ -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<HandlerResult> => {
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,
};
}
};
};
Loading