diff --git a/apps/gateway/src/app.ts b/apps/gateway/src/app.ts index 0641c4e..496c879 100644 --- a/apps/gateway/src/app.ts +++ b/apps/gateway/src/app.ts @@ -930,6 +930,8 @@ function assertQuoteStillExecutable( original: z.infer, refreshed: z.infer, ) { + const originalSourceBlock = BigInt(original.sourceBlock); + const refreshedSourceBlock = BigInt(refreshed.sourceBlock); const fields = [ "network", "marketId", @@ -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.", diff --git a/apps/gateway/test/trading.test.mjs b/apps/gateway/test/trading.test.mjs index 572b95a..0174e81 100644 --- a/apps/gateway/test/trading.test.mjs +++ b/apps/gateway/test/trading.test.mjs @@ -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", @@ -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; @@ -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(); }); diff --git a/render.yaml b/render.yaml index c478b67..59da360 100644 --- a/render.yaml +++ b/render.yaml @@ -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 diff --git a/scripts/production-smoke.mjs b/scripts/production-smoke.mjs index 5a7860d..d293569 100644 --- a/scripts/production-smoke.mjs +++ b/scripts/production-smoke.mjs @@ -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, diff --git a/scripts/test/production-smoke.test.mjs b/scripts/test/production-smoke.test.mjs index 7cadcd5..d605979 100644 --- a/scripts/test/production-smoke.test.mjs +++ b/scripts/test/production-smoke.test.mjs @@ -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); @@ -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", @@ -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"); }); diff --git a/workers/market-sync/src/runner.ts b/workers/market-sync/src/runner.ts index 0deafef..ad1551f 100644 --- a/workers/market-sync/src/runner.ts +++ b/workers/market-sync/src/runner.ts @@ -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()); @@ -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 { return new Promise((resolve) => { const timer = setTimeout(resolve, milliseconds);