Skip to content
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