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
227 changes: 227 additions & 0 deletions indexer/streams/src/db/repository.ts
Original file line number Diff line number Diff line change
@@ -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<void>;
fundStream(streamId: string, amount: string, eventKey?: string): Promise<void>;
recordWithdrawal(
streamId: string,
recipient: string,
amount: string,
txHash: string,
timestamp: string,
eventKey?: string,
): Promise<void>;
recordCancel(
streamId: string,
canceler: string,
txHash: string,
timestamp: string,
eventKey?: string,
): Promise<void>;
}

export class StreamWriteServiceImpl implements StreamWriteService {
private processedEventKeys = new Set<string>();

constructor(private readonly repository: StreamRepository) {}

async createStream(input: StreamCreateInput, eventKey?: string): Promise<void> {
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<void> {
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<void> {
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<void> {
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<Stream>;
private withdrawalRepo: Repository<WithdrawalAction>;
private cancelRepo: Repository<CancelAction>;

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<void> {
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<Stream | null> {
return this.streamRepo.findOneBy({ id: streamId });
}

async fundStream(streamId: string, amount: string): Promise<void> {
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<void> {
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<void> {
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 });
});
}
}
70 changes: 43 additions & 27 deletions indexer/streams/src/handlers/stream-cancel.handler.ts
Original file line number Diff line number Diff line change
@@ -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<HandlerResult> => {
try {
const payload = parseStreamCancel(event.data);
export function createStreamCancelHandler(persistence?: StreamWriteService): 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.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();
86 changes: 48 additions & 38 deletions indexer/streams/src/handlers/stream-funded.handler.ts
Original file line number Diff line number Diff line change
@@ -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<HandlerResult> => {
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<HandlerResult> => {
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();
Loading