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
Original file line number Diff line number Diff line change
@@ -1,7 +1,10 @@
import { describe, expect, it } from "bun:test";
import { RampDirection } from "@vortexfi/shared";
import { Op } from "sequelize";
import { RAMP_START_EXPIRATION_TIME_SECONDS } from "../../../../../constants/constants";
import { getPersistedBlockFlowCompatibilityScope } from "./compatibility-scope";
import { getFundedInitialSellRampWhere, getPersistedBlockFlowCompatibilityScope } from "./compatibility-scope";

const THREE_DAYS_MS = 3 * 24 * 60 * 60 * 1000;

describe("persisted block-flow compatibility scope", () => {
it("scopes pending quotes and resumable ramps to the current flow variant", () => {
Expand All @@ -19,9 +22,28 @@ describe("persisted block-flow compatibility scope", () => {
[Op.or]: [
{ currentPhase: { [Op.notIn]: ["complete", "failed", "timedOut", "initial"] } },
{ createdAt: { [Op.gte]: initialRampCutoff }, currentPhase: "initial" },
{ currentPhase: "initial", "state.aveniaTicketId": { [Op.ne]: null } }
{ currentPhase: "initial", "state.aveniaTicketId": { [Op.ne]: null } },
// Funded SELL ramps the recovery worker may start past the client window.
{
createdAt: { [Op.gt]: new Date(now.getTime() - THREE_DAYS_MS), [Op.lt]: now },
currentPhase: "initial",
type: RampDirection.SELL,
[Op.or]: [
{ "state.squidRouterSwapHash": { [Op.ne]: null } },
{ "state.squidRouterNoPermitTransferHash": { [Op.ne]: null } }
]
}
]
}
});
});

it("selects funded SELL ramps only between the minimum age and the three-day recovery window", () => {
const now = new Date("2026-07-31T12:00:00.000Z");

expect(getFundedInitialSellRampWhere(now, 16 * 60 * 1000).createdAt).toEqual({
[Op.gt]: new Date(now.getTime() - THREE_DAYS_MS),
[Op.lt]: new Date(now.getTime() - 16 * 60 * 1000)
});
});
});
Original file line number Diff line number Diff line change
@@ -1,17 +1,43 @@
import { RampDirection } from "@vortexfi/shared";
import { Op } from "sequelize";
import type { FlowVariant } from "../../../../../config/vars";
import { RAMP_START_EXPIRATION_TIME_SECONDS } from "../../../../../constants/constants";

const TERMINAL_RAMP_PHASES = ["complete", "failed", "timedOut"] as const;
const FUNDED_SELL_RECOVERY_WINDOW_MS = 3 * 24 * 60 * 60 * 1000;

/**
* `initial` SELL ramps whose user-broadcast source transaction hash was already reported, created
* within the recovery window and at least `minAgeMs` ago. The user's funds are on the ephemeral
* once that transaction mines, so the recovery worker starts these past the client start window.
* The worker selects with this predicate and the startup check keeps their flow versions
* registered, so the two cannot drift apart.
*/
export function getFundedInitialSellRampWhere(now = new Date(), minAgeMs = 0) {
return {
createdAt: {
[Op.gt]: new Date(now.getTime() - FUNDED_SELL_RECOVERY_WINDOW_MS),
[Op.lt]: new Date(now.getTime() - minAgeMs)
},
currentPhase: "initial" as const,
type: RampDirection.SELL,
[Op.or]: [
{ "state.squidRouterSwapHash": { [Op.ne]: null } },
{ "state.squidRouterNoPermitTransferHash": { [Op.ne]: null } }
]
};
}

/**
* Selects only persisted state that this backend could still execute.
*
* A registered ramp remains in `initial` until startRamp is called. Both updateRamp
* and the public startRamp reject it after the shared expiration window. Avenia ramps
* are the exception: registration creates a payable PIX ticket, and the recovery
* worker may start an expired initial ramp after the provider confirms payment. Those
* rows therefore remain deployment dependencies. Once a ramp has entered a financial
* and the public startRamp reject it after the shared expiration window. Two kinds of
* ramp are the exception: Avenia registration creates a payable PIX ticket, and the
* recovery worker may start an expired initial ramp after the provider confirms payment;
* and a SELL ramp whose user already reported its source transaction hash is started by
* the same worker (see getFundedInitialSellRampWhere). Those rows therefore remain
* deployment dependencies. Once a ramp has entered a financial
* phase, age never makes it safe to ignore: every non-terminal phase owned by this flow
* variant stays fail-closed.
*/
Expand All @@ -29,7 +55,8 @@ export function getPersistedBlockFlowCompatibilityScope(flowVariant: FlowVariant
[Op.or]: [
{ currentPhase: { [Op.notIn]: [...TERMINAL_RAMP_PHASES, "initial"] } },
{ createdAt: { [Op.gte]: initialRampCutoff }, currentPhase: "initial" },
{ currentPhase: "initial", "state.aveniaTicketId": { [Op.ne]: null } }
{ currentPhase: "initial", "state.aveniaTicketId": { [Op.ne]: null } },
getFundedInitialSellRampWhere(now)
]
}
};
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -71,7 +71,7 @@ describe("RampService Moonbeam retirement", () => {
expect(update).not.toHaveBeenCalled();
});

it("rejects public and provider-paid starts before persisted flow execution", async () => {
it("rejects public, provider-paid, and funded-SELL starts before persisted flow execution", async () => {
RampState.findByPk = mock(async () => ({
createdAt: new Date(),
currentPhase: "initial",
Expand All @@ -85,7 +85,11 @@ describe("RampService Moonbeam retirement", () => {
})) as unknown as typeof RampState.findByPk;

const service = new TestRampService();
for (const start of [() => service.startRamp({ rampId: "ramp-1" }), () => service.recoverPaidAveniaRamp("ramp-1")]) {
for (const start of [
() => service.startRamp({ rampId: "ramp-1" }),
() => service.recoverPaidAveniaRamp("ramp-1"),
() => service.recoverFundedSellRamp("ramp-1")
]) {
await expect(start()).rejects.toMatchObject({ status: httpStatus.SERVICE_UNAVAILABLE });
}
});
Expand Down
Original file line number Diff line number Diff line change
@@ -0,0 +1,81 @@
import { afterEach, describe, expect, it, mock } from "bun:test";
import { FiatToken, Networks, RampDirection } from "@vortexfi/shared";
import httpStatus from "http-status";
import type { Transaction } from "sequelize";
import { config } from "../../../config/vars";
import QuoteTicket from "../../../models/quoteTicket.model";
import RampState from "../../../models/rampState.model";
import { RampService } from "./ramp.service";

class TestRampService extends RampService {
protected async withTransaction<T>(callback: (transaction: Transaction) => Promise<T>): Promise<T> {
return callback({} as Transaction);
}
}

const originalQuoteFindByPk = QuoteTicket.findByPk;
const originalRampFindByPk = RampState.findByPk;

afterEach(() => {
QuoteTicket.findByPk = originalQuoteFindByPk;
RampState.findByPk = originalRampFindByPk;
});

function stubRampAndQuote(
ramp: { from?: Networks; state: Record<string, unknown>; to?: string; type: RampDirection },
outputCurrency: string
) {
RampState.findByPk = mock(async () => ({
createdAt: new Date(Date.now() - 60 * 60 * 1000),
currentPhase: "initial",
flowVariant: config.flowVariant,
from: Networks.Ethereum,
id: "ramp-1",
presignedTxs: [],
quoteId: "quote-1",
to: "pix",
unsignedTxs: [],
...ramp
})) as unknown as typeof RampState.findByPk;
QuoteTicket.findByPk = mock(async () => ({
id: "quote-1",
metadata: { blocks: {}, flow: { id: "BrlOfframpBase" }, globals: { fees: { usd: {} }, request: {} } },
outputCurrency
})) as unknown as typeof QuoteTicket.findByPk;
}

describe("RampService.recoverFundedSellRamp guards", () => {
const conflict = { message: "Ramp does not have a reported source transaction", status: httpStatus.CONFLICT };

it("refuses a SELL ramp whose source transaction hash was never reported", async () => {
stubRampAndQuote({ state: {}, type: RampDirection.SELL }, FiatToken.BRL);

await expect(new TestRampService().recoverFundedSellRamp("ramp-1")).rejects.toMatchObject(conflict);
});

it("refuses a BUY ramp even when a hash-shaped field is present", async () => {
stubRampAndQuote({ state: { squidRouterSwapHash: "0xabc" }, type: RampDirection.BUY }, FiatToken.BRL);

await expect(new TestRampService().recoverFundedSellRamp("ramp-1")).rejects.toMatchObject(conflict);
});

it("refuses a domestic (AlfredPay) SELL whose reported hash FundEphemeral does not verify", async () => {
stubRampAndQuote({ state: { squidRouterNoPermitTransferHash: "0xabc" }, type: RampDirection.SELL }, FiatToken.MXN);

await expect(new TestRampService().recoverFundedSellRamp("ramp-1")).rejects.toMatchObject(conflict);
});

it("refuses an AssetHub SELL whose reported Squid hash FundEphemeral does not verify", async () => {
stubRampAndQuote(
{
from: Networks.AssetHub,
state: { assethubToPendulumHash: "0xdef", squidRouterSwapHash: "0xabc" },
to: "sepa",
type: RampDirection.SELL
},
FiatToken.EURC
);

await expect(new TestRampService().recoverFundedSellRamp("ramp-1")).rejects.toMatchObject(conflict);
});
});
31 changes: 30 additions & 1 deletion apps/api/src/api/services/ramp/ramp.service.ts
Original file line number Diff line number Diff line change
Expand Up @@ -576,9 +576,22 @@ export class RampService extends BaseRampService {
return this.startRampWithOptions({ rampId }, { enforceDeadline: false, requirePaidAveniaTicket: true });
}

/**
* Start an EVM SELL ramp whose user already reported the hash of their source transaction but
* whose client never reached /ramp/start inside the window. That transaction delivers the funds
* to the ephemeral, so the deadline no longer protects anyone; FundEphemeral verifies the
* reported hash against the issued blueprint on-chain before any platform spend.
*/
public async recoverFundedSellRamp(rampId: string): Promise<StartRampResponse> {
return this.startRampWithOptions(
{ rampId },
{ enforceDeadline: false, requirePaidAveniaTicket: false, requireReportedSellSource: true }
);
}

private async startRampWithOptions(
request: StartRampRequest,
options: { enforceDeadline: boolean; requirePaidAveniaTicket: boolean }
options: { enforceDeadline: boolean; requirePaidAveniaTicket: boolean; requireReportedSellSource?: boolean }
): Promise<StartRampResponse> {
return this.withTransaction(async transaction => {
const rampState = await RampState.findByPk(request.rampId, { lock: Transaction.LOCK.UPDATE, transaction });
Expand Down Expand Up @@ -622,6 +635,22 @@ export class RampService extends BaseRampService {
status: httpStatus.CONFLICT
});
}
if (options.requireReportedSellSource) {
// Domestic (AlfredPay) and AssetHub SELLs are excluded: FundEphemeral only verifies the
// reported hash for the other EVM SELLs, so this recovery has no pre-spend proof for them.
const { squidRouterNoPermitTransferHash, squidRouterSwapHash } = rampState.state;
if (
rampState.type !== RampDirection.SELL ||
rampState.from === Networks.AssetHub ||
isDomesticToken(quote.outputCurrency as FiatToken) ||
!(squidRouterSwapHash || squidRouterNoPermitTransferHash)
Comment on lines +643 to +646
) {
throw new APIError({
message: "Ramp does not have a reported source transaction",
status: httpStatus.CONFLICT
});
}
}
if (options.enforceDeadline) {
RampService.assertStartDeadlineNotExceeded(rampState);
}
Expand Down
94 changes: 92 additions & 2 deletions apps/api/src/api/workers/ramp-recovery.worker.test.ts
Original file line number Diff line number Diff line change
@@ -1,7 +1,11 @@
import { afterEach, beforeEach, describe, expect, it, mock } from "bun:test";
import { EPaymentMethod, Networks } from "@vortexfi/shared";
import { afterEach, beforeEach, describe, expect, it, mock, spyOn } from "bun:test";
import { EPaymentMethod, Networks, RampDirection } from "@vortexfi/shared";
import { Op } from "sequelize";
import logger from "../../config/logger";
import { config } from "../../config/vars";
import RampState from "../../models/rampState.model";
import phaseProcessor from "../services/phases/phase-processor";
import rampService from "../services/ramp/ramp.service";
import RampRecoveryWorker from "./ramp-recovery.worker";

const originalFindAll = RampState.findAll;
Expand Down Expand Up @@ -37,3 +41,89 @@ describe("RampRecoveryWorker Moonbeam retirement", () => {
expect(processRamp).not.toHaveBeenCalled();
});
});

describe("RampRecoveryWorker funded SELL start", () => {
const originalRecoverFundedSellRamp = rampService.recoverFundedSellRamp;
const originalAppendErrorLog = rampService.appendErrorLog;
const recoverFundedSellRamp = mock(async (_rampId: string): Promise<unknown> => undefined);
const appendErrorLog = mock(async (_id: string, _entry: unknown) => undefined);
const fundedSell = {
currentPhase: "initial",
from: Networks.Ethereum,
id: "funded-sell-ramp",
state: { flow: { id: "BrlOfframpBase" }, squidRouterSwapHash: "0xabc" },
to: EPaymentMethod.PIX,
unsignedTxs: []
};
let queries: Array<{ where: Record<PropertyKey, unknown> }>;

beforeEach(() => {
queries = [];
RampState.findAll = mock(async (options: { where: Record<PropertyKey, unknown> }) => {
queries.push(options);
return options.where.currentPhase === "initial" ? [fundedSell] : [];
}) as unknown as typeof RampState.findAll;
rampService.recoverFundedSellRamp = recoverFundedSellRamp as unknown as typeof rampService.recoverFundedSellRamp;
rampService.appendErrorLog = appendErrorLog as unknown as typeof rampService.appendErrorLog;
recoverFundedSellRamp.mockReset();
recoverFundedSellRamp.mockImplementation(async () => undefined);
appendErrorLog.mockClear();
});

afterEach(() => {
rampService.recoverFundedSellRamp = originalRecoverFundedSellRamp;
rampService.appendErrorLog = originalAppendErrorLog;
});

async function runWorker() {
const worker = new RampRecoveryWorker("*/5 * * * *", false) as unknown as { recover: () => Promise<void> };
await worker.recover();
}

it("selects initial SELL ramps with a reported source hash between 16 minutes and 3 days old", async () => {
const before = Date.now();
await runWorker();
const after = Date.now();

const where = queries.find(query => query.where.currentPhase === "initial")?.where as Record<PropertyKey, unknown>;
expect(where.type).toBe(RampDirection.SELL);
expect(where.flowVariant).toBe(config.flowVariant);
expect(where[Op.or]).toEqual([
{ "state.squidRouterSwapHash": { [Op.ne]: null } },
{ "state.squidRouterNoPermitTransferHash": { [Op.ne]: null } }
]);
const createdAt = where.createdAt as Record<symbol, Date>;
const minute = 60 * 1000;
// The worker reads the clock between `before` and `after`, so each cutoff lies in that window.
expect(createdAt[Op.lt].getTime()).toBeGreaterThanOrEqual(before - 16 * minute);
expect(createdAt[Op.lt].getTime()).toBeLessThanOrEqual(after - 16 * minute);
expect(createdAt[Op.gt].getTime()).toBeGreaterThanOrEqual(before - 3 * 24 * 60 * minute);
expect(createdAt[Op.gt].getTime()).toBeLessThanOrEqual(after - 3 * 24 * 60 * minute);
});

it("starts each selected ramp through the funded SELL path, not the phase processor", async () => {
await runWorker();

expect(recoverFundedSellRamp).toHaveBeenCalledTimes(1);
expect(recoverFundedSellRamp).toHaveBeenCalledWith("funded-sell-ramp");
expect(processRamp).not.toHaveBeenCalled();
expect(appendErrorLog).not.toHaveBeenCalled();
});

it("logs a failed start on the ramp and selects it again on the next cycle", async () => {
recoverFundedSellRamp.mockImplementation(async () => {
throw new Error("database unavailable");
});

const info = spyOn(logger, "info");
await runWorker();
await runWorker();

expect(info).toHaveBeenCalledWith("Ramp recovery attempt completed. Successful: 0, Failed: 1");
info.mockRestore();
expect(appendErrorLog).toHaveBeenCalledTimes(2);
expect(appendErrorLog.mock.calls[0]?.[0]).toBe("funded-sell-ramp");
expect(appendErrorLog.mock.calls[0]?.[1]).toMatchObject({ error: "database unavailable", phase: "initial" });
expect(recoverFundedSellRamp).toHaveBeenCalledTimes(2);
});
});
Loading
Loading