Skip to content
Merged
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
9 changes: 7 additions & 2 deletions apps/gateway/src/app.ts
Original file line number Diff line number Diff line change
Expand Up @@ -930,6 +930,8 @@ function assertQuoteStillExecutable(
original: z.infer<typeof TradeQuoteSchema>,
refreshed: z.infer<typeof TradeQuoteSchema>,
) {
const originalSourceBlock = BigInt(original.sourceBlock);
const refreshedSourceBlock = BigInt(refreshed.sourceBlock);
const fields = [
"network",
"marketId",
Expand All @@ -947,9 +949,12 @@ function assertQuoteStillExecutable(
"limitPrice",
"tickSize",
"lotSize",
"sourceBlock",
] as const;
if (fields.some((field) => original[field] !== refreshed[field])) {
if (
refreshedSourceBlock < originalSourceBlock ||
refreshedSourceBlock - originalSourceBlock > 20n ||
fields.some((field) => original[field] !== refreshed[field])
) {
throw new QuoteError(
"STALE_BOOK",
"The executable quote no longer matches authoritative DreamDEX depth.",
Expand Down
50 changes: 48 additions & 2 deletions apps/gateway/test/trading.test.mjs
Original file line number Diff line number Diff line change
Expand Up @@ -135,7 +135,18 @@ test("empty books return an actionable conflict instead of a fabricated quote",
});

test("plan endpoint binds the account and current immutable market generation", async () => {
const quoteApp = createGateway({ dataReader: reader() });
const advancingReader = reader();
const getOrderBook = advancingReader.getOrderBook;
let sourceBlock = 92n;
advancingReader.getOrderBook = async (...args) => {
const book = await getOrderBook(...args);
sourceBlock += 8n;
return {
...book,
freshness: { ...book.freshness, sourceBlock: sourceBlock.toString() },
};
};
const quoteApp = createGateway({ dataReader: advancingReader });
const quoteResponse = await quoteApp.inject({
method: "POST",
url: "/v1/trading/quotes",
Expand All @@ -145,7 +156,7 @@ test("plan endpoint binds the account and current immutable market generation",
await quoteApp.close();
let planned;
const app = createGateway({
dataReader: reader(),
dataReader: advancingReader,
tradePlanner: {
create: async (input) => {
planned = input;
Expand Down Expand Up @@ -173,6 +184,41 @@ test("plan endpoint binds the account and current immutable market generation",
assert.equal(planned.chainId, 50312);
assert.equal(planned.account, account);
assert.equal(planned.quote.marketId, marketId);
assert.equal(planned.quote.sourceBlock, "100");
await app.close();
});

test("plan endpoint rejects a quote when refreshed depth changes its executable terms", async () => {
const quoteApp = createGateway({ dataReader: reader() });
const quoteResponse = await quoteApp.inject({
method: "POST",
url: "/v1/trading/quotes",
payload: { marketId, outcome: "up", side: "buy", mode: "quantity", amount: "100000" },
});
const quote = quoteResponse.json();
await quoteApp.close();
const app = createGateway({
dataReader: reader([{ price: "700000", quantity: "10000000" }]),
tradePlanner: { create: async () => ({ accepted: true }) },
});
const response = await app.inject({
method: "POST",
url: "/v1/trading/plans",
payload: {
account,
idempotencyKey: "gateway-plan-changed-depth",
quote,
policy: {
maxSlippageBps: 100,
minimumFillBps: 10000,
minimumTimeRemainingSeconds: 30,
quoteSourceBlock: quote.sourceBlock,
quoteExpiresAt: quote.expiresAt,
},
},
});
assert.equal(response.statusCode, 409);
assert.equal(response.json().error.code, "STALE_BOOK");
await app.close();
});

Expand Down
4 changes: 3 additions & 1 deletion render.yaml
Original file line number Diff line number Diff line change
Expand Up @@ -39,7 +39,9 @@ services:
- key: CORS_ALLOWED_ORIGINS
value: https://eventrail.vercel.app
- key: WORKER_INTERVAL_MS
value: "15000"
value: "30000"
- key: MARKET_SYNC_HISTORICAL_LIMIT
value: "50"
# The free preview runs only the live market synchronizer. Receipt,
# settlement, and webhook workers belong on separate production services.
- key: DEMO_WORKERS
Expand Down
5 changes: 4 additions & 1 deletion scripts/production-smoke.mjs
Original file line number Diff line number Diff line change
Expand Up @@ -36,7 +36,10 @@ await check("claims", new URL(`/v1/data/accounts/${account}/claims`, gatewayUrl)
auth: true,
});
await checkStream(new URL("/v1/events?network=shannon", gatewayUrl));
const market = Array.isArray(markets) ? markets[0] : markets?.data?.[0];
const marketRows = Array.isArray(markets) ? markets : (markets?.data ?? []);
const market = [...marketRows]
.filter((candidate) => candidate?.marketId)
.sort((left, right) => Date.parse(right.expiresAt ?? 0) - Date.parse(left.expiresAt ?? 0))[0];
if (market?.marketId) {
const quote = await post("quote", "/v1/trading/quotes", {
marketId: market.marketId,
Expand Down
13 changes: 10 additions & 3 deletions scripts/test/production-smoke.test.mjs
Original file line number Diff line number Diff line change
Expand Up @@ -9,6 +9,7 @@ const smokeScript = fileURLToPath(new URL("../production-smoke.mjs", import.meta
test("production smoke supports separate application and gateway origins", async (context) => {
const applicationRequests = [];
const gatewayRequests = [];
const quoteRequests = [];
const planRequests = [];
const application = createServer((request, response) => {
applicationRequests.push(request.url);
Expand All @@ -30,16 +31,21 @@ test("production smoke supports separate application and gateway origins", async
response.end(": connected\n\n");
return;
}
if (request.url === "/v1/trading/plans") {
if (request.url === "/v1/trading/quotes" || request.url === "/v1/trading/plans") {
let body = "";
request.setEncoding("utf8");
for await (const chunk of request) body += chunk;
planRequests.push(JSON.parse(body));
const payload = JSON.parse(body);
if (request.url === "/v1/trading/quotes") quoteRequests.push(payload);
else planRequests.push(payload);
}
response.writeHead(200, { "content-type": "application/json" });
response.end(
request.url?.startsWith("/v1/data/markets?")
? JSON.stringify([{ marketId: `0x${"ab".repeat(32)}` }])
? JSON.stringify([
{ marketId: `0x${"ab".repeat(32)}`, expiresAt: "2029-01-01T00:00:00.000Z" },
{ marketId: `0x${"cd".repeat(32)}`, expiresAt: "2031-01-01T00:00:00.000Z" },
])
: request.url === "/v1/trading/quotes"
? JSON.stringify({
quoteId: "test",
Expand Down Expand Up @@ -70,6 +76,7 @@ test("production smoke supports separate application and gateway origins", async
assert.ok(gatewayRequests.some((request) => request.startsWith("GET /v1/data/markets?")));
assert.ok(gatewayRequests.includes("POST /v1/trading/quotes"));
assert.ok(gatewayRequests.includes("POST /v1/trading/plans"));
assert.equal(quoteRequests[0].marketId, `0x${"cd".repeat(32)}`);
assert.equal(planRequests[0].policy.quoteSourceBlock, "123456");
assert.equal(planRequests[0].policy.quoteExpiresAt, "2030-01-01T00:00:00.000Z");
});
Expand Down
8 changes: 8 additions & 0 deletions workers/market-sync/src/runner.ts
Original file line number Diff line number Diff line change
Expand Up @@ -12,6 +12,7 @@ const synchronizer = new MarketSynchronizer(
createDreamDexAdapter({ network: "shannon" }),
new PostgresDreamDexMarketGenerationRepository(database),
new DreamDexEventBus(redis),
{ historicalLimit: positiveInteger(process.env.MARKET_SYNC_HISTORICAL_LIMIT ?? "250") },
);

for (const signal of ["SIGINT", "SIGTERM"] as const) process.once(signal, () => controller.abort());
Expand All @@ -34,6 +35,13 @@ function safeError(error: unknown) {
: { message: "Unknown failure" };
}

function positiveInteger(value: string): number {
const parsed = Number(value);
if (!/^[1-9][0-9]*$/.test(value) || !Number.isSafeInteger(parsed))
throw new Error("MARKET_SYNC_HISTORICAL_LIMIT must be a positive integer");
return parsed;
}

function delay(milliseconds: number, signal: AbortSignal): Promise<void> {
return new Promise((resolve) => {
const timer = setTimeout(resolve, milliseconds);
Expand Down
Loading