From 7d622821e8e5a149fedfbd883d53c71665f690bc Mon Sep 17 00:00:00 2001 From: my908-hue Date: Wed, 30 Sep 2026 13:01:14 +0000 Subject: [PATCH] feat(intents): add real-time RFQ quote auctions --- .env.example | 3 + .env.mainnet.example | 1 + .env.staging.example | 1 + .env.testnet.example | 1 + docs/solver-onboarding.md | 32 ++++ scripts/db-migrate-locked.spec.ts | 2 +- src/common/rfq-signature.ts | 15 ++ src/common/stellar-signature.ts | 3 - src/config/configuration.ts | 15 +- src/config/env.validation.ts | 13 +- src/intents/dto/quote-response.dto.ts | 5 +- src/intents/feed/intent-feed.service.ts | 137 ++++++++++++++- src/intents/intents.controller.ts | 179 ++++++++++++++------ src/intents/intents.gateway.ts | 141 ++++----------- src/intents/intents.service.ts | 4 +- src/intents/rfq.types.ts | 35 ++++ src/soroban/fill-verifier.service.ts | 36 +++- src/soroban/redaction.ts | 2 +- src/soroban/signer.service.spec.ts | 28 --- src/soroban/solver-registry.service.spec.ts | 22 +-- src/soroban/soroban.service.ts | 25 +-- src/treasury/treasury.service.ts | 5 +- test/chaos/runner.ts | 4 + 23 files changed, 443 insertions(+), 266 deletions(-) create mode 100644 src/common/rfq-signature.ts create mode 100644 src/intents/rfq.types.ts diff --git a/.env.example b/.env.example index 00b5cbe7..8bb9169e 100644 --- a/.env.example +++ b/.env.example @@ -96,6 +96,9 @@ SOROBAN_FEE_PERCENTILE=p50 INTENT_RETENTION_DAYS=30 INTENT_RETENTION_SWEEP_MS=60000 +# Time allowed to collect connected solver RFQ responses (1-1000 ms). +QUOTE_AUCTION_WINDOW_MS=300 + # ─── CORS ──────────────────────────────────────────────────────────────────── # Comma-separated list of allowed origins for the frontend. # Development default: "*" (any origin allowed — convenient for local work) diff --git a/.env.mainnet.example b/.env.mainnet.example index 1564f9f3..9cfcb470 100644 --- a/.env.mainnet.example +++ b/.env.mainnet.example @@ -77,6 +77,7 @@ CORS_ORIGIN=https://app.vortex.trade # # ─── WebSocket ─────────────────────────────────────────────────────────────── # Tune based on expected solver + frontend connection count. WS_MAX_CONNECTIONS=5000 +QUOTE_AUCTION_WINDOW_MS=300 # ─── Pluggable signer backend (issue #400) ─────────────────────────────────── # REQUIRED in production: use SIGNER_BACKEND=vault so the signing key never diff --git a/.env.staging.example b/.env.staging.example index 293a442c..9bed51d5 100644 --- a/.env.staging.example +++ b/.env.staging.example @@ -34,6 +34,7 @@ SOROBAN_FEE_PERCENTILE=p50 CORS_ORIGIN=* WS_MAX_CONNECTIONS=1000 +QUOTE_AUCTION_WINDOW_MS=300 # ─── Pluggable signer backend (issue #400) ─────────────────────────────────── SIGNER_BACKEND=local diff --git a/.env.testnet.example b/.env.testnet.example index 38d41076..0be61ebc 100644 --- a/.env.testnet.example +++ b/.env.testnet.example @@ -68,6 +68,7 @@ CORS_ORIGIN=* # ─── WebSocket ─────────────────────────────────────────────────────────────── WS_MAX_CONNECTIONS=1000 +QUOTE_AUCTION_WINDOW_MS=300 # ─── Pluggable signer backend (issue #400) ─────────────────────────────────── # SIGNER_BACKEND=local is the default for development. diff --git a/docs/solver-onboarding.md b/docs/solver-onboarding.md index 18c57779..62e90e84 100644 --- a/docs/solver-onboarding.md +++ b/docs/solver-onboarding.md @@ -153,6 +153,38 @@ Upon subscription, the WebSocket server responds with a `subscribed` event: ``` Subsequent `intent_created` events will only be broadcast to the bot if the intent's `srcChain` matches one of the subscribed chains. +### RFQ Quote Requests +Authenticated, active solvers with a positive bond and matching source-chain/token capabilities may receive a short-lived `rfq_request` over this WebSocket. The default response window is 300 ms and can be configured from 1 to 1,000 ms with `QUOTE_AUCTION_WINDOW_MS`. + +```json +{ + "type": "rfq_request", + "requestId": "550e8400-e29b-41d4-a716-446655440000", + "srcChain": "ethereum", + "srcTokenSymbol": "USDC", + "srcAmount": "1000000", + "dstTokenSymbol": "USDC", + "deadline": 1775836800300 +} +``` + +Reply before `deadline` with the gross destination amount, solver fee, expiry in Unix seconds, and a Stellar Ed25519 signature: + +```json +{ + "type": "rfq_response", + "requestId": "550e8400-e29b-41d4-a716-446655440000", + "dstAmount": "998500", + "fee": "100", + "expiresAt": 1775836860, + "signature": "base64EncodedSignatureString==" +} +``` + +Sign the UTF-8 bytes of `vortex:rfq:v1::`. `payloadHash` is the lowercase SHA-256 hex digest of the JSON encoding of the request fields (`requestId`, `srcChain`, `srcTokenSymbol`, `srcAmount`, `dstTokenSymbol`, optional `srcTokenAddress` and `dstTokenContract`, and `deadline`) plus `solver`, `dstAmount`, `fee`, and `expiresAt`. Omit absent optional fields and sort keys lexicographically before `JSON.stringify`. The signature is Base64-encoded. The backend accepts one valid response per solver and request; late, expired, malformed, or invalidly signed responses are ignored. + +Quotes are ranked by destination amount after solver and protocol fees, with reputation breaking ties. If no valid solver response arrives within the window, the API returns the existing model-based estimate with `indicative: true`; otherwise `indicative` is false. + ### Event Replay & Reconnection On connection or reconnection, the bot can request event replay from its last received sequence ID (`seq`) to avoid missing intents during network blips: ```json diff --git a/scripts/db-migrate-locked.spec.ts b/scripts/db-migrate-locked.spec.ts index 4efc0145..14cb54ea 100644 --- a/scripts/db-migrate-locked.spec.ts +++ b/scripts/db-migrate-locked.spec.ts @@ -13,7 +13,7 @@ * `require` and typed here instead of via an ES import (no allowJs in * tsconfig, and the runtime must not depend on generated types). */ -// eslint-disable-next-line @typescript-eslint/no-require-imports +// eslint-disable-next-line @typescript-eslint/no-var-requires const migrate = require("./db-migrate-locked.js") as { CHECKPOINT_DDL: string; CHECKPOINT_TABLE: string; diff --git a/src/common/rfq-signature.ts b/src/common/rfq-signature.ts new file mode 100644 index 00000000..2cad43d7 --- /dev/null +++ b/src/common/rfq-signature.ts @@ -0,0 +1,15 @@ +import { createHash } from "node:crypto"; +import { RfqResponseSignaturePayload } from "../intents/rfq.types"; + +/** Canonical domain-separated message signed by a solver for an RFQ response. */ +export function buildRfqResponseMessage(payload: RfqResponseSignaturePayload): string { + const canonicalPayload = JSON.stringify( + Object.fromEntries( + Object.entries(payload) + .filter(([, value]) => value !== undefined) + .sort(([left], [right]) => left.localeCompare(right)), + ), + ); + const payloadHash = createHash("sha256").update(canonicalPayload, "utf8").digest("hex"); + return `vortex:rfq:v1:${payload.requestId}:${payloadHash}`; +} \ No newline at end of file diff --git a/src/common/stellar-signature.ts b/src/common/stellar-signature.ts index f51e7a80..a8623f61 100644 --- a/src/common/stellar-signature.ts +++ b/src/common/stellar-signature.ts @@ -31,9 +31,6 @@ function canonicalPayload(payload: Record): stri function buildV2IntentMessage( context: IntentSignatureContext, action: "accept" | "fill" | "cancel", - /** - * Build the canonical message that a solver must sign to update their mutable - * profile fields (name / supportedChains / supportedTokens / avgFillTime). intentId: string, payload: Record, ): string { diff --git a/src/config/configuration.ts b/src/config/configuration.ts index cf883776..2035f2c2 100644 --- a/src/config/configuration.ts +++ b/src/config/configuration.ts @@ -99,7 +99,6 @@ export interface AppConfig { nodeEnv: string; port: number; databaseUrl: string; - datasets: import("../datasets/datasets.types").DatasetsConfig; stellar: { network: "testnet" | "futurenet" | "mainnet"; sorobanRpcUrl: string; @@ -129,6 +128,7 @@ export interface AppConfig { }; intentRetentionDays: number; intentRetentionSweepMs: number; + quoteAuctionWindowMs: number; /** * Dry-run flag for on-chain write paths (issue #260). * @@ -272,6 +272,7 @@ export interface AppConfig { refreshIntervalMs: number; /** Comma-separated extra secrets: "name:envVar:required". */ extra: string; + }; /** WS gateway hardening (issue #455). */ ws: { /** Largest inbound frame accepted; larger frames close the socket (1009). */ @@ -370,6 +371,7 @@ export default (): AppConfig => ({ }, intentRetentionDays: parseInt(process.env.INTENT_RETENTION_DAYS ?? "30", 10), intentRetentionSweepMs: parseInt(process.env.INTENT_RETENTION_SWEEP_MS ?? "60000", 10), + quoteAuctionWindowMs: parseInt(process.env.QUOTE_AUCTION_WINDOW_MS ?? "300", 10), // Default to dry-run (true) outside production; in production the value must // be explicitly set (validated by envValidationSchema). onchainDryRun: process.env.ONCHAIN_DRY_RUN !== undefined @@ -429,16 +431,6 @@ export default (): AppConfig => ({ overrides: process.env.FLAG_OVERRIDES ?? "", }, adminApiKeys: process.env.ADMIN_API_KEYS ?? "", - datasets: { - enabled: (process.env.DATASETS_ENABLED ?? "false") === "true", - anonymize: (process.env.DATASETS_ANONYMIZE ?? "true") === "true", - salt: process.env.DATASETS_SALT ?? "", - saltRotationHours: parseInt(process.env.DATASETS_SALT_ROTATION_HOURS ?? "24", 10), - saltRetentionWindows: parseInt(process.env.DATASETS_SALT_RETENTION_WINDOWS ?? "2", 10), - publicBucket: process.env.DATASETS_PUBLIC_BUCKET ?? "", - storageKind: (process.env.DATASETS_STORAGE_KIND ?? "memory") as "local" | "memory", - localDir: process.env.DATASETS_LOCAL_DIR ?? "", - }, guardianContractId: process.env.GUARDIAN_CONTRACT_ID ?? "", canaryAddresses: (process.env.CANARY_ADDRESSES ?? "") .split(",") @@ -458,6 +450,7 @@ export default (): AppConfig => ({ provider: (process.env.SECRETS_PROVIDER ?? "env") as "env" | "aws-secrets-manager" | "vault-kv", refreshIntervalMs: parseInt(process.env.SECRETS_REFRESH_INTERVAL_MS ?? "60000", 10), extra: process.env.SECRETS_EXTRA ?? "", + }, ws: { maxPayloadBytes: parseInt(process.env.WS_MAX_PAYLOAD_BYTES ?? "16384", 10), maxConnectionsPerIp: parseInt(process.env.WS_MAX_CONNECTIONS_PER_IP ?? "20", 10), diff --git a/src/config/env.validation.ts b/src/config/env.validation.ts index 300a9416..ff499169 100644 --- a/src/config/env.validation.ts +++ b/src/config/env.validation.ts @@ -50,9 +50,6 @@ export const envValidationSchema = Joi.object({ }), ONCHAIN_INTENTS_ENABLED: Joi.boolean().default(false), - // Stellar public key of the treasury account (fee/slash/refund accumulator). - TREASURY_ADDRESS: Joi.string().allow("").default(""), - // Stellar public key of the treasury account (fee accumulator). TREASURY_ADDRESS: Joi.string().allow("").default(""), @@ -107,6 +104,7 @@ export const envValidationSchema = Joi.object({ // the eviction sweep runs. Both are read by IntentsService. INTENT_RETENTION_DAYS: Joi.number().integer().min(0).default(30), INTENT_RETENTION_SWEEP_MS: Joi.number().integer().min(0).default(60000), + QUOTE_AUCTION_WINDOW_MS: Joi.number().integer().min(1).max(1000).default(300), // ── Reference solver bot (scripts/solver-bot.ts) ─────────────────────────── // Read by the standalone bot process rather than by the server, but declared @@ -353,15 +351,8 @@ export const envValidationSchema = Joi.object({ .pattern(/^([A-Za-z0-9_.-]+:(admin|superadmin):[^,:]{16,})(,[A-Za-z0-9_.-]+:(admin|superadmin):[^,:]{16,})*$/) .default(""), - // ── Public anonymised datasets ──────────────────────────────────────────── - DATASETS_ENABLED: Joi.boolean().default(false), - DATASETS_ANONYMIZE: Joi.boolean().default(true), - DATASETS_SALT: Joi.string().allow("").default(""), - DATASETS_SALT_ROTATION_HOURS: Joi.number().integer().min(1).max(720).default(24), - DATASETS_SALT_RETENTION_WINDOWS: Joi.number().integer().min(0).max(30).default(2), - DATASETS_PUBLIC_BUCKET: Joi.string().allow("").default(""), + // Legacy storage selector retained for existing deployments. DATASETS_STORAGE_KIND: Joi.string().valid("local", "memory").default("memory"), - DATASETS_LOCAL_DIR: Joi.string().allow("").default(""), // ── Guardian emergency ingestion (issue #507) ───────────────────────────── GUARDIAN_CONTRACT_ID: Joi.string().allow("").default(""), diff --git a/src/intents/dto/quote-response.dto.ts b/src/intents/dto/quote-response.dto.ts index ee9ddcf7..cbbe944f 100644 --- a/src/intents/dto/quote-response.dto.ts +++ b/src/intents/dto/quote-response.dto.ts @@ -1,4 +1,4 @@ -import { ApiProperty } from "@nestjs/swagger"; +import { ApiProperty, ApiPropertyOptional } from "@nestjs/swagger"; import { TokenInfo } from "../intents.types"; export class RouteStepDto { @@ -100,4 +100,7 @@ export class QuoteResponseDto { @ApiProperty({ description: "Price impact for the best quote as a decimal fraction (0 when no quote available)" }) priceImpact!: number; + + @ApiPropertyOptional({ description: "True when no solver responded and the returned quote is indicative" }) + indicative?: boolean; } diff --git a/src/intents/feed/intent-feed.service.ts b/src/intents/feed/intent-feed.service.ts index aa08f00c..2b5c1a52 100644 --- a/src/intents/feed/intent-feed.service.ts +++ b/src/intents/feed/intent-feed.service.ts @@ -13,6 +13,17 @@ import { MemoryBackplane } from "../backplane/memory.backplane"; import { EventRingBuffer } from "../event-ring-buffer"; import { FeedClient, FeedAdmission, FeedFilter, FeedReplayResult } from "./feed.types"; import { resolveClientIp } from "../ws/connection-state"; +import { verifyStellarSignature } from "../../common/stellar-signature"; +import { buildRfqResponseMessage } from "../../common/rfq-signature"; +import { RfqQuoteRequest, RfqResponseSignaturePayload, VerifiedRfqQuote } from "../rfq.types"; + +interface PendingRfq { + request: RfqQuoteRequest; + eligibleSolvers: Set; + responses: Map; + resolve: (responses: VerifiedRfqQuote[]) => void; + timer: ReturnType; +} /** * How many sequenced events to keep in the replay buffer. @@ -46,6 +57,7 @@ export class IntentFeedService implements OnModuleDestroy { /** Connected clients and their per-connection filters. */ private readonly clients = new Map(); + private readonly pendingRfqs = new Map(); /** Per-IP connection accounting (shared by WS and SSE). */ private readonly connectionsPerIp = new Map(); @@ -150,6 +162,119 @@ export class IntentFeedService implements OnModuleDestroy { return this.backplane.publish(event); } + requestRfq( + request: Omit, + eligibleSolvers: string[], + windowMs: number, + ): Promise { + const eligible = new Set(eligibleSolvers); + if (eligible.size === 0) return Promise.resolve([]); + + const rfq: RfqQuoteRequest = { + ...request, + requestId: randomUUID(), + deadline: Date.now() + windowMs, + }; + + return new Promise((resolve) => { + const timer = setTimeout(() => this.completeRfq(rfq.requestId), windowMs); + this.pendingRfqs.set(rfq.requestId, { + request: rfq, + eligibleSolvers: eligible, + responses: new Map(), + resolve, + timer, + }); + + void this.backplane + .publish({ type: "rfq_request", ...rfq, eligibleSolvers: [...eligible] }) + .catch(() => this.completeRfq(rfq.requestId)); + }); + } + + submitRfqResponse(response: { + solver: string; + requestId: string; + dstAmount: string; + fee: string; + expiresAt: number; + signature: string; + }): Promise { + return this.backplane.publish({ type: "rfq_response", ...response }); + } + + private completeRfq(requestId: string): void { + const pending = this.pendingRfqs.get(requestId); + if (!pending) return; + clearTimeout(pending.timer); + this.pendingRfqs.delete(requestId); + pending.resolve([...pending.responses.values()]); + } + + private deliverRfqRequest(event: SequencedEvent): void { + const eligibleSolvers = (event as SequencedEvent & { eligibleSolvers?: unknown }).eligibleSolvers; + if (!Array.isArray(eligibleSolvers)) return; + const eligible = new Set(eligibleSolvers.filter((solver): solver is string => typeof solver === "string")); + const { seq, eligibleSolvers: _eligibleSolvers, ...request } = event as SequencedEvent & { + eligibleSolvers: string[]; + }; + const payload = JSON.stringify({ seq, ...request }); + + for (const [client, filter] of this.clients) { + if (filter.solver && eligible.has(filter.solver.solverAddress)) { + this.sendToClient(client, payload, seq); + } + } + } + + private acceptRfqResponse(event: SequencedEvent): void { + const response = event as SequencedEvent & { + solver?: unknown; + requestId?: unknown; + dstAmount?: unknown; + fee?: unknown; + expiresAt?: unknown; + signature?: unknown; + }; + const { solver, requestId, dstAmount, fee, expiresAt, signature } = response; + if ( + typeof solver !== "string" || + typeof requestId !== "string" || + typeof dstAmount !== "string" || + typeof fee !== "string" || + typeof expiresAt !== "number" || + typeof signature !== "string" + ) return; + + const pending = this.pendingRfqs.get(requestId); + if ( + !pending || + Date.now() >= pending.request.deadline || + !pending.eligibleSolvers.has(solver) || + pending.responses.has(solver) + ) return; + if (!/^\d{1,78}$/.test(dstAmount) || !/^\d{1,78}$/.test(fee) || !Number.isSafeInteger(expiresAt)) return; + + const now = Math.floor(Date.now() / 1000); + if (expiresAt <= now || expiresAt > now + 60) return; + + try { + if (BigInt(dstAmount) <= BigInt(fee)) return; + const payload: RfqResponseSignaturePayload = { + ...pending.request, + solver, + dstAmount, + fee, + expiresAt, + }; + verifyStellarSignature(solver, buildRfqResponseMessage(payload), signature); + } catch { + return; + } + + pending.responses.set(solver, { solver, dstAmount, fee, expiresAt }); + } + /** Chains deliveries so async chain lookups cannot reorder events. */ private enqueueDelivery(event: SequencedEvent): Promise { const run = this.deliveryChain.then(() => this.deliver(event)); @@ -167,6 +292,15 @@ export class IntentFeedService implements OnModuleDestroy { const enqueuedAt = Date.now(); const { seq, ...event } = sequencedEvent; + if (event.type === "rfq_request") { + this.deliverRfqRequest(sequencedEvent); + return; + } + if (event.type === "rfq_response") { + this.acceptRfqResponse(sequencedEvent); + return; + } + this.updateIndexForEvent(sequencedEvent); this.ringBuffer.push(sequencedEvent); @@ -283,7 +417,7 @@ export class IntentFeedService implements OnModuleDestroy { /** Current sequence number (0 when no events have been broadcast). */ get currentSeq(): number { - return this.ringBuffer.latestSeq(); + return this.backplane.health().lastSeq; } // ── Solver capability ─────────────────────────────────────────────────── @@ -367,6 +501,7 @@ export class IntentFeedService implements OnModuleDestroy { } async onModuleDestroy(): Promise { + for (const requestId of this.pendingRfqs.keys()) this.completeRfq(requestId); await this.backplane.close(); for (const [client] of this.clients) { this.removeClient(client); diff --git a/src/intents/intents.controller.ts b/src/intents/intents.controller.ts index 3be5fc33..fe3d86d4 100644 --- a/src/intents/intents.controller.ts +++ b/src/intents/intents.controller.ts @@ -32,7 +32,7 @@ import { import { Throttle } from "@nestjs/throttler"; import { IntentsService } from "./intents.service"; import { IntentsGateway } from "./intents.gateway"; -import { SolversService } from "../solvers/solvers.service"; +import { SolversService, solverSupports } from "../solvers/solvers.service"; import { TokensService } from "../tokens/tokens.service"; import { RoutingService } from "../routing/routing.service"; import { MAX_OPEN_INTENTS_PER_USER } from "./intents.service"; @@ -60,6 +60,7 @@ import { } from "../common/stellar-signature"; import { SignatureNonceService } from "../common/signature-nonce.service"; import { EvmSignatureVerifier } from "../common/evm-signature"; +import { SolverRecord } from "../solvers/solvers.types"; import { applyVarianceScale, calculateProtocolFee, @@ -90,6 +91,16 @@ interface IntentSignatureProof { expiresAt?: number; } +interface QuoteCandidate { + solver: SolverRecord; + dstAmount: bigint; + fee: bigint; + netDstAmount: bigint; + fillTime: number; + expiresAt: number; + reputationScore: number; +} + @ApiTags("intents") @Controller("api/v1/intents") export class IntentsController { @@ -111,6 +122,7 @@ export class IntentsController { this.signatureNetwork = config.get("stellar.network", { infer: true }); this.legacyStellarSignatures = config.get("legacyStellarSignatures", { infer: true }); this.nodeEnv = config.get("nodeEnv", { infer: true }); + this.quoteAuctionWindowMs = config.get("quoteAuctionWindowMs", { infer: true }); } /** Canary addresses (issue #496). */ @@ -118,6 +130,7 @@ export class IntentsController { private readonly signatureNetwork: AppConfig["stellar"]["network"]; private readonly legacyStellarSignatures: boolean; private readonly nodeEnv: string; + private readonly quoteAuctionWindowMs: number; private verifyIntentSignature( action: IntentSignatureProof["action"], @@ -853,7 +866,12 @@ export class IntentsController { }) @ApiOkResponse({ type: QuoteResponseDto }) async quote(@Body() dto: QuoteRequestDto): Promise { - const solvers = (await this.solversService.getAll()).filter((s) => s.isActive); + const solvers = (await this.solversService.getAll()).filter( + (solver) => + solver.isActive && + BigInt(solver.bondAmount) > 0n && + solverSupports(solver, dto.srcChain, dto.srcTokenSymbol), + ); // #219: use typed resolveSrcToken / resolveDstToken — no more any casts. // #276: a quote may be requested by symbol alone (no contract/address), but @@ -875,8 +893,24 @@ export class IntentsController { // eslint-disable-next-line @typescript-eslint/no-explicit-any const srcPriceUSD: number = (srcToken as any)?.priceUSD ?? dstPriceUSD; - const quotes = solvers - .map((solver) => { + const rfqResponses = await this.intentsGateway.requestRfq( + { + srcChain: dto.srcChain, + srcTokenSymbol: dto.srcTokenSymbol, + srcAmount: dto.srcAmount, + dstTokenSymbol: dto.dstTokenSymbol, + srcTokenAddress: dto.srcTokenAddress, + dstTokenContract: dto.dstTokenContract, + }, + solvers.map((solver) => solver.address), + this.quoteAuctionWindowMs, + ); + const indicative = rfqResponses.length === 0; + const now = Math.floor(Date.now() / 1000); + const quoteCandidates: QuoteCandidate[] = []; + + if (indicative) { + for (const solver of solvers) { // Issue #118: weight variance by solver performance history. const totalFills = solver.fillsCompleted + solver.fillsFailed; const successRate = totalFills > 0 ? solver.fillsCompleted / totalFills : 0.5; @@ -885,59 +919,95 @@ export class IntentsController { const varianceScaled = varianceScaleFromPerfScore(perfScore); const dstAmount = applyVarianceScale(srcAmountBigInt, varianceScaled); const fee = calculateProtocolFee(dstAmount); // 0.05% + const ageDays = Math.max(0, (now - solver.registeredAt) / 86_400); + const reputationScore = successRate * Math.exp(-ageDays / 180); + quoteCandidates.push({ + solver, + dstAmount, + fee, + netDstAmount: dstAmount > fee ? dstAmount - fee : 0n, + fillTime: solver.avgFillTime + Math.floor(Math.random() * 30), + expiresAt: now + 60, + reputationScore, + }); + } + } else { + const solversByAddress = new Map(solvers.map((solver) => [solver.address, solver])); + for (const response of rfqResponses) { + const solver = solversByAddress.get(response.solver); + if (!solver) continue; + const dstAmount = BigInt(response.dstAmount); + const fee = BigInt(response.fee) + calculateProtocolFee(dstAmount); + if (fee >= dstAmount) continue; + const totalFills = solver.fillsCompleted + solver.fillsFailed; + const successRate = totalFills > 0 ? solver.fillsCompleted / totalFills : 0; + const ageDays = Math.max(0, (now - solver.registeredAt) / 86_400); + quoteCandidates.push({ + solver, + dstAmount, + fee, + netDstAmount: dstAmount - fee, + fillTime: solver.avgFillTime, + expiresAt: response.expiresAt, + reputationScore: successRate * Math.exp(-ageDays / 180), + }); + } + } - // Issue #126: compute USD fee total and price impact. - // eslint-disable-next-line @typescript-eslint/no-explicit-any - const feeUnits = toDecimalNumber(fee, (dstToken as any)?.decimals ?? 7); - const totalFeesUSD = feeUnits * dstPriceUSD; - const srcUnits = toDecimalNumber(srcAmountBigInt, srcToken?.decimals ?? 7); - const dstUnits = toDecimalNumber(dstAmount, dstToken?.decimals ?? 7); - const priceImpact = - srcPriceUSD > 0 && dstPriceUSD > 0 - ? Math.max(0, 1 - (dstUnits * dstPriceUSD) / (srcUnits * srcPriceUSD)) - : 0; - - // #220: attach a computed route to each solver quote. - // Build minimal TokenInfo objects for routing (uses resolved data when available). - const srcTokenInfo = { - address: dto.srcTokenAddress ?? "", - symbol: dto.srcTokenSymbol, - name: srcToken?.name ?? dto.srcTokenSymbol, - decimals: srcToken?.decimals ?? 18, - chain: (dto.srcChain as SupportedChain) ?? "ethereum", - priceUSD: srcToken?.priceUSD, - }; - const dstTokenInfo = { - address: dstToken?.contract ?? dto.dstTokenContract ?? "", - symbol: dto.dstTokenSymbol, - name: dstToken?.name ?? dto.dstTokenSymbol, - decimals: dstToken?.decimals ?? 7, - chain: "stellar" as SupportedChain, - priceUSD: dstToken?.priceUSD, - }; + quoteCandidates.sort((left, right) => { + if (left.netDstAmount !== right.netDstAmount) { + return left.netDstAmount > right.netDstAmount ? -1 : 1; + } + return right.reputationScore - left.reputationScore || right.solver.fillsCompleted - left.solver.fillsCompleted; + }); - // Try a direct route; fall back to a two-hop via USDC intermediate when - // a direct solver path is not viable (different base tokens). - const route = this.routingService.buildRoute(srcTokenInfo, dstTokenInfo, solver.address, { - totalFeesUSD, - priceImpact, - estimatedFillTime: solver.avgFillTime + Math.floor(Math.random() * 30), - }); + const quotes = quoteCandidates.map((candidate) => { + // Issue #126: compute USD fee total and price impact. + // eslint-disable-next-line @typescript-eslint/no-explicit-any + const feeUnits = toDecimalNumber(candidate.fee, (dstToken as any)?.decimals ?? 7); + const totalFeesUSD = feeUnits * dstPriceUSD; + const srcUnits = toDecimalNumber(srcAmountBigInt, srcToken?.decimals ?? 7); + const netDstUnits = toDecimalNumber(candidate.netDstAmount, dstToken?.decimals ?? 7); + const priceImpact = + srcPriceUSD > 0 && dstPriceUSD > 0 + ? Math.max(0, 1 - (netDstUnits * dstPriceUSD) / (srcUnits * srcPriceUSD)) + : 0; + + const srcTokenInfo = { + address: dto.srcTokenAddress ?? "", + symbol: dto.srcTokenSymbol, + name: srcToken?.name ?? dto.srcTokenSymbol, + decimals: srcToken?.decimals ?? 18, + chain: dto.srcChain, + priceUSD: srcToken?.priceUSD, + }; + const dstTokenInfo = { + address: dstToken?.contract ?? dto.dstTokenContract ?? "", + symbol: dto.dstTokenSymbol, + name: dstToken?.name ?? dto.dstTokenSymbol, + decimals: dstToken?.decimals ?? 7, + chain: "stellar" as SupportedChain, + priceUSD: dstToken?.priceUSD, + }; + const route = this.routingService.buildRoute( + srcTokenInfo, + dstTokenInfo, + candidate.solver.address, + { totalFeesUSD, priceImpact, estimatedFillTime: candidate.fillTime }, + ); - return { - solver: solver.address, - solverName: solver.name, - dstAmount: dstAmount.toString(), - fee: fee.toString(), - fillTime: solver.avgFillTime + Math.floor(Math.random() * 30), - expiresAt: Math.floor(Date.now() / 1000) + 60, - totalFeesUSD, - priceImpact, - route, - }; - }) - // nosemgrep: no-number-money -- sort comparator on bounded quote diffs only; amounts stay strings elsewhere. - .sort((a, b) => Number(BigInt(b.dstAmount) - BigInt(a.dstAmount))); + return { + solver: candidate.solver.address, + solverName: candidate.solver.name, + dstAmount: candidate.dstAmount.toString(), + fee: candidate.fee.toString(), + fillTime: candidate.fillTime, + expiresAt: candidate.expiresAt, + totalFeesUSD, + priceImpact, + route, + }; + }); if (dto.intentId && quotes.length > 0) { await this.intentsService.update(dto.intentId, { quotedDstAmount: quotes[0].dstAmount }); @@ -954,6 +1024,7 @@ export class IntentsController { estimatedFillTime: best?.fillTime ?? 0, totalFeesUSD: best?.totalFeesUSD ?? 0, priceImpact: best?.priceImpact ?? 0, + indicative, }; } diff --git a/src/intents/intents.gateway.ts b/src/intents/intents.gateway.ts index 1f66f448..727fa0a4 100644 --- a/src/intents/intents.gateway.ts +++ b/src/intents/intents.gateway.ts @@ -24,6 +24,7 @@ import { ConnectionState, resolveClientIp } from "./ws/connection-state"; import { IntentFeedService } from "./feed/intent-feed.service"; import { FeedClient, FeedFilter } from "./feed/feed.types"; import { randomUUID } from "node:crypto"; +import { RfqQuoteRequest, VerifiedRfqQuote } from "./rfq.types"; export type { SequencedEvent } from "./backplane/backplane.types"; export { EventRingBuffer } from "./event-ring-buffer"; @@ -91,114 +92,12 @@ export class IntentsGateway return this.feed.backplaneHealth(); } - /** - * Broadcast an event to every connected client (WS and SSE) on every replica. - * - * Activity 1: Uses encoding cache — serializes once per format, not per client. - * - * Delivery rules (evaluated in order): - * 1. Client is not OPEN → skip. - * 2. Client set wantAll=true → always deliver. - * 3. Client has a solver capability predicate: - * a. Event carries an inlined intent → apply predicate to that intent. - * b. Event is a state-transition (only intentId available) → deliver - * (we cannot efficiently look up the intent here; the solver would - * already have received the intent_created event through the filter). - * 4. Client has a plain chain filter (`chains != null`) → apply chain match. - * 5. No filter → full unfiltered feed (backward-compatible default). - */ - private deliverToMatchingSubscribers( - seq: number, - chain: SupportedChain | null, - event: { type: string; [key: string]: unknown }, - ) { - for (const [client, filter] of this.subscribers) { - if (client.readyState !== WebSocket.OPEN) continue; - - // Opt-out: solver requested full feed. - if (filter.wantAll) { - this.sendEncoded(client, seq, event); - continue; - } - - // Authenticated solver — apply capability predicate. - if (filter.solver !== null) { - const solverPredicate = filter.solver; - const inlinedIntent = (event as { intent?: unknown }).intent; - - // intent_created carries a full intent object we can test directly. - if (event.type === "intent_created" && inlinedIntent && typeof inlinedIntent === "object") { - // eslint-disable-next-line @typescript-eslint/no-explicit-any - const matches = solverPredicate.matches(inlinedIntent as any); - if (matches) { - this.sendEncoded(client, seq, event); - try { this.metricsService?.incWsDelivered(solverPredicate.solverAddress); } catch { /* noop */ } - } else { - try { this.metricsService?.incWsFiltered(solverPredicate.solverAddress); } catch { /* noop */ } - } - continue; - } - - // State-transition events: the solver already filtered on intent_created, - // so we pass them through to keep the feed self-consistent. - this.sendEncoded(client, seq, event); - try { this.metricsService?.incWsDelivered(solverPredicate.solverAddress); } catch { /* noop */ } - continue; - } - - // No filter set → full unfiltered feed (backward-compatible default). - if (filter.chains === null) { - this.sendEncoded(client, seq, event); - continue; - } - - // Chain couldn't be resolved → deliver to everyone (safe default). - if (chain === null) { - this.sendEncoded(client, seq, event); - continue; - } - - // Only send if the event's chain is in this subscriber's filter. - if (filter.chains.has(chain)) { - this.sendEncoded(client, seq, event); - } - } - } - - handleConnection(client: WebSocket) { - this.subscribers.set(client, { - chains: null, - solver: null, - wantAll: false, - subscriptionCount: 0, - }); - /** - * Send an event to a client using its negotiated encoding format (Activity 1). - * - * Retrieves pre-serialized payload from encoding cache, avoiding redundant - * serialization work. At 10k connections with msgpack, this saves 9,999 - * msgpackEncode() calls per broadcast. - */ - private sendEncoded(client: WebSocket, seq: number, event: Record): void { - const state = this.connections.get(client); - if (!state) { - // Fallback for connections without state (shouldn't happen) - if (client.readyState === WebSocket.OPEN) { - client.send(JSON.stringify({ seq, ...event })); - } - return; - } - - const payload = this.encodingCache.get(seq, event, state.encoding); - const result = state.send(payload); - - if (result === "dropped_oldest") { - this.metricsService?.wsOutboundDropped.inc(); - } else if (result === "disconnected") { - this.metricsService?.wsSlowConsumerDisconnects.inc(); - logger.warn(`ws slow consumer disconnected (ip=${state.ip}, queue full)`); - this.removeSubscriber(client); - } + requestRfq( + request: Omit, + eligibleSolvers: string[], + windowMs: number, + ): Promise { + return this.feed.requestRfq(request, eligibleSolvers, windowMs); } /** @@ -441,11 +340,37 @@ export class IntentsGateway case "auth": await this.handleAuth(client, msg); break; + case "rfq_response": + this.handleRfqResponse(client, msg); + break; default: break; } } + private handleRfqResponse(client: WebSocket, payload: Record): void { + const solver = this.authenticatedSolver.get(client); + if ( + !solver || + typeof payload.requestId !== "string" || + payload.requestId.length > 64 || + typeof payload.dstAmount !== "string" || + typeof payload.fee !== "string" || + typeof payload.expiresAt !== "number" || + typeof payload.signature !== "string" || + payload.signature.length > 128 + ) return; + + void this.feed.submitRfqResponse({ + solver, + requestId: payload.requestId, + dstAmount: payload.dstAmount, + fee: payload.fee, + expiresAt: payload.expiresAt, + signature: payload.signature, + }).catch(() => undefined); + } + /** * Process a `{ type: "subscribe", chains?: string[], all?: boolean }` message. * diff --git a/src/intents/intents.service.ts b/src/intents/intents.service.ts index c54cd91d..9032b37e 100644 --- a/src/intents/intents.service.ts +++ b/src/intents/intents.service.ts @@ -706,8 +706,6 @@ export class IntentsService { // the sweep loop already logs that case loudly. const slashedSolver = subject?.solver; if (subject && slashedSolver) { - const solver = subject?.solver; - if (solver && subject) { this.reportShadow( "slash", subject.intentId, @@ -716,7 +714,7 @@ export class IntentsService { this.safeArgs(() => [ nativeToScVal(subject.intentId, { type: "string" }), new Address(slashedSolver).toScVal(), - new Address(solver).toScVal(), + new Address(slashedSolver).toScVal(), nativeToScVal(patch.slashReason, { type: "string" }), nativeToScVal(patch.slashedAt, { type: "u64" }), ]), diff --git a/src/intents/rfq.types.ts b/src/intents/rfq.types.ts new file mode 100644 index 00000000..2941252d --- /dev/null +++ b/src/intents/rfq.types.ts @@ -0,0 +1,35 @@ +import { SupportedChain } from "./intents.types"; + +export interface RfqQuoteRequest { + requestId: string; + srcChain: SupportedChain; + srcTokenSymbol: string; + srcAmount: string; + dstTokenSymbol: string; + srcTokenAddress?: string; + dstTokenContract?: string; + deadline: number; +} + +export interface RfqResponseSignaturePayload extends RfqQuoteRequest { + solver: string; + dstAmount: string; + fee: string; + expiresAt: number; +} + +export interface RfqQuoteResponse { + type: "rfq_response"; + requestId: string; + dstAmount: string; + fee: string; + expiresAt: number; + signature: string; +} + +export interface VerifiedRfqQuote { + solver: string; + dstAmount: string; + fee: string; + expiresAt: number; +} \ No newline at end of file diff --git a/src/soroban/fill-verifier.service.ts b/src/soroban/fill-verifier.service.ts index c29041ca..716469c3 100644 --- a/src/soroban/fill-verifier.service.ts +++ b/src/soroban/fill-verifier.service.ts @@ -1,6 +1,7 @@ import { Injectable } from "@nestjs/common"; import { ConfigService } from "@nestjs/config"; import { Asset } from "@stellar/stellar-sdk"; +import { EgressPurpose, HttpEgressService } from "../common/http-egress"; import { AppConfig, NETWORK_PASSPHRASES } from "../config/configuration"; import { Intent } from "../intents/intents.types"; @@ -20,20 +21,35 @@ type HorizonTransaction = { /** Independently checks Horizon's indexed Stellar transaction and payment operations. */ @Injectable() export class FillVerifierService { - constructor(private readonly config: ConfigService) {} + private readonly egress: HttpEgressService; + private readonly horizonBase: string; + + constructor(private readonly config: ConfigService) { + const horizonUrl = config.get("stellar.horizonUrl", { infer: true }); + this.horizonBase = horizonUrl.replace(/\/$/, ""); + this.egress = new HttpEgressService({ + timeoutMs: 10_000, + maxRedirects: 0, + maxBodySizeBytes: 1_048_576, + allowlist: [new URL(horizonUrl).hostname], + blockPrivateRanges: false, + }); + } /** * Verify the transaction against the persisted intent. Unknown/indexing errors * remain retryable; malformed or mismatched transactions are definitive. */ async verify(txHash: string, intent: Intent): Promise { - const base = this.config.get("stellar.horizonUrl", { infer: true }).replace(/\/$/, ""); + const base = this.horizonBase; let transaction: HorizonTransaction; try { - const response = await fetch(`${base}/transactions/${encodeURIComponent(txHash)}`); - if (response.status === 404) return { status: "pending", reason: "not_indexed" }; - if (!response.ok) return { status: "pending", reason: "horizon_unavailable" }; - transaction = (await response.json()) as HorizonTransaction; + const response = await this.egress.fetch(`${base}/transactions/${encodeURIComponent(txHash)}`, { + purpose: EgressPurpose.HORIZON, + }); + if (response.statusCode === 404) return { status: "pending", reason: "not_indexed" }; + if (response.statusCode < 200 || response.statusCode >= 300) return { status: "pending", reason: "horizon_unavailable" }; + transaction = JSON.parse(response.body) as HorizonTransaction; } catch { return { status: "pending", reason: "horizon_unavailable" }; } @@ -44,9 +60,11 @@ export class FillVerifierService { } if (!transaction._links?.operations?.href) return { status: "rejected", reason: "operations_missing" }; try { - const response = await fetch(`${base}/transactions/${encodeURIComponent(txHash)}/operations?limit=200&order=asc`); - if (!response.ok) return { status: "pending", reason: "horizon_unavailable" }; - const body = (await response.json()) as { _embedded?: { records?: Array> } }; + const response = await this.egress.fetch(`${base}/transactions/${encodeURIComponent(txHash)}/operations?limit=200&order=asc`, { + purpose: EgressPurpose.HORIZON, + }); + if (response.statusCode < 200 || response.statusCode >= 300) return { status: "pending", reason: "horizon_unavailable" }; + const body = JSON.parse(response.body) as { _embedded?: { records?: Array> } }; const operations = body._embedded?.records ?? []; for (const operation of operations) { const type = operation.type as string; diff --git a/src/soroban/redaction.ts b/src/soroban/redaction.ts index 36d681fa..769b9424 100644 --- a/src/soroban/redaction.ts +++ b/src/soroban/redaction.ts @@ -18,7 +18,7 @@ const SENSITIVE_KEY_PATTERNS = [ // AWS secret access key /AKIA[0-9A-Z]{16}/g, // Generic password/secret in URL - /:\/\/[^:\/\s]+:([^@\/\s]{8,})@/gi, + /:\/\/[^:/\s]+:([^@/\s]{8,})@/gi, // Database connection strings with passwords /postgresql:\/\/[^:]+:([^@]+)@/gi, /mysql:\/\/[^:]+:([^@]+)@/gi, diff --git a/src/soroban/signer.service.spec.ts b/src/soroban/signer.service.spec.ts index 80df543d..d13afd70 100644 --- a/src/soroban/signer.service.spec.ts +++ b/src/soroban/signer.service.spec.ts @@ -4,20 +4,6 @@ import { AppConfig } from "../config/configuration"; import { SignerService } from "./signer.service"; import { SorobanService } from "./soroban.service"; import { ISigner } from "./signers/signer.interface"; -import { LocalKeypairSigner } from "./signers/local-keypair.signer"; -import { findSensitiveKeyMaterial } from "./redaction"; - -/** - * Builds a local-keypair signing backend over a stubbed config, i.e. the same - * ISigner the SIGNER_TOKEN provider hands to SignerService in the app. - */ -function signerWith(signerSecretKey: string, network: AppConfig["stellar"]["network"] = "testnet"): ISigner { - const values: Record = { - "stellar.signingKey": signerSecretKey, - "stellar.network": network, - }; - const configService = { get: (path: string) => values[path] } as ConfigService; - return new LocalKeypairSigner(configService); import { findSensitiveKeyMaterial } from "./redaction"; /** @@ -53,14 +39,11 @@ function fakeSorobanService(startingSequence = "100") { describe("SignerService", () => { it("reports unconfigured when no secret is set", () => { - const service = new SignerService(signerWith(""), fakeSorobanService()); const service = new SignerService(fakeSigner(), fakeSorobanService()); expect(service.isConfigured()).toBe(false); }); it("throws a clear, secret-free error when signing without a configured key", () => { - const service = new SignerService(signerWith(""), fakeSorobanService()); - expect(() => service.getPublicKey()).toThrow(/SOROBAN_SIGNING_KEY/); const service = new SignerService(fakeSigner(), fakeSorobanService()); expect(() => service.getPublicKey()).toThrow(/not configured|signing key/i); // The message must stay secret-free (no strkey / seed shape). @@ -75,7 +58,6 @@ describe("SignerService", () => { it("derives the public key from the configured secret", () => { const keypair = Keypair.random(); - const service = new SignerService(signerWith(keypair.secret()), fakeSorobanService()); const service = new SignerService(fakeSigner(keypair), fakeSorobanService()); expect(service.isConfigured()).toBe(true); @@ -84,9 +66,6 @@ describe("SignerService", () => { it("maps network config to the right passphrase", () => { const soroban = fakeSorobanService(); - expect(new SignerService(signerWith("", "testnet"), soroban).getNetworkPassphrase()).toBe(Networks.TESTNET); - expect(new SignerService(signerWith("", "futurenet"), soroban).getNetworkPassphrase()).toBe(Networks.FUTURENET); - expect(new SignerService(signerWith("", "mainnet"), soroban).getNetworkPassphrase()).toBe(Networks.PUBLIC); expect(new SignerService(fakeSigner(undefined, "testnet"), soroban).getNetworkPassphrase()).toBe(Networks.TESTNET); expect(new SignerService(fakeSigner(undefined, "futurenet"), soroban).getNetworkPassphrase()).toBe(Networks.FUTURENET); expect(new SignerService(fakeSigner(undefined, "mainnet"), soroban).getNetworkPassphrase()).toBe(Networks.PUBLIC); @@ -94,7 +73,6 @@ describe("SignerService", () => { it("signs a transaction with the configured key", async () => { const keypair = Keypair.random(); - const service = new SignerService(signerWith(keypair.secret()), fakeSorobanService()); const service = new SignerService(fakeSigner(keypair), fakeSorobanService()); const account = new Account(keypair.publicKey(), "1"); @@ -110,7 +88,6 @@ describe("SignerService", () => { it("never includes the raw secret in string/JSON/inspect representations", () => { const keypair = Keypair.random(); - const service = new SignerService(signerWith(keypair.secret()), fakeSorobanService()); const service = new SignerService(fakeSigner(keypair), fakeSorobanService()); const secret = keypair.secret(); @@ -136,7 +113,6 @@ describe("SignerService", () => { it("fetches the starting sequence once and increments it locally", async () => { const keypair = Keypair.random(); const soroban = fakeSorobanService("100"); - const service = new SignerService(signerWith(keypair.secret()), soroban); const service = new SignerService(fakeSigner(keypair), soroban); const first = await service.withNextSequence(async (sequence) => sequence); @@ -150,7 +126,6 @@ describe("SignerService", () => { it("hands out a distinct, gap-free sequence to every concurrent caller", async () => { const keypair = Keypair.random(); const soroban = fakeSorobanService("0"); - const service = new SignerService(signerWith(keypair.secret()), soroban); const service = new SignerService(fakeSigner(keypair), soroban); const results = await Promise.all( @@ -164,7 +139,6 @@ describe("SignerService", () => { it("runs callers strictly one at a time, in call order", async () => { const keypair = Keypair.random(); - const service = new SignerService(signerWith(keypair.secret()), fakeSorobanService("0")); const service = new SignerService(fakeSigner(keypair), fakeSorobanService("0")); const order: number[] = []; @@ -183,7 +157,6 @@ describe("SignerService", () => { it("drops the cached sequence after a failure so the next call re-syncs from the network", async () => { const keypair = Keypair.random(); const soroban = fakeSorobanService("100"); - const service = new SignerService(signerWith(keypair.secret()), soroban); const service = new SignerService(fakeSigner(keypair), soroban); await expect( @@ -199,7 +172,6 @@ describe("SignerService", () => { it("does not let a failed caller block callers queued behind it", async () => { const keypair = Keypair.random(); - const service = new SignerService(signerWith(keypair.secret()), fakeSorobanService("0")); const service = new SignerService(fakeSigner(keypair), fakeSorobanService("0")); const failing = service.withNextSequence(async () => { diff --git a/src/soroban/solver-registry.service.spec.ts b/src/soroban/solver-registry.service.spec.ts index 88134db3..fe356ada 100644 --- a/src/soroban/solver-registry.service.spec.ts +++ b/src/soroban/solver-registry.service.spec.ts @@ -81,6 +81,7 @@ function makeConfigService( provider: "env", refreshIntervalMs: 60000, extra: "", + }, ws: { maxPayloadBytes: 16384, maxConnectionsPerIp: 20, @@ -91,6 +92,7 @@ function makeConfigService( outboundQueueMax: 1000, outboundBufferBytes: 1048576, slowConsumerPolicy: "drop_oldest", + drainTimeoutMs: 25000, }, authJwtSecret: "", rateLimitLocalPruneMs: 60000, @@ -100,26 +102,6 @@ function makeConfigService( heartbeatMs: 15000, maxBufferBytes: 1048576, }, - datasets: { - enabled: false, - anonymize: true, - salt: "", - saltRotationHours: 24, - saltRetentionWindows: 2, - publicBucket: "", - storageKind: "memory", - localDir: "", - }, - treasury: { - address: "", - }, - shadow: { - enabled: false, - sampleRate: 1, - queueMax: 256, - concurrency: 4, - sourceAccount: "", - }, health: { roles: ["api", "ws", "worker"], checkIntervalMs: 5000, diff --git a/src/soroban/soroban.service.ts b/src/soroban/soroban.service.ts index f22558a9..4e90af05 100644 --- a/src/soroban/soroban.service.ts +++ b/src/soroban/soroban.service.ts @@ -1,17 +1,26 @@ import { Injectable } from "@nestjs/common"; import { ConfigService } from "@nestjs/config"; import { SorobanRpc, Transaction } from "@stellar/stellar-sdk"; +import { EgressPurpose, HttpEgressService } from "../common/http-egress"; import { AppConfig } from "../config/configuration"; @Injectable() export class SorobanService { private readonly server: SorobanRpc.Server; private readonly rpcUrl: string; + private readonly egress: HttpEgressService; constructor(configService: ConfigService) { const rpcUrl = configService.get("stellar.sorobanRpcUrl", { infer: true }); this.rpcUrl = rpcUrl; this.server = new SorobanRpc.Server(rpcUrl, { allowHttp: rpcUrl.startsWith("http://") }); + this.egress = new HttpEgressService({ + timeoutMs: 10_000, + maxRedirects: 0, + maxBodySizeBytes: 1_048_576, + allowlist: [new URL(rpcUrl).hostname], + blockPrivateRanges: false, + }); } getHealth() { @@ -42,7 +51,7 @@ export class SorobanService { * the `{ header: { closeTime } }` shape the ingestion loop reads. */ async getLedger(sequence: number): Promise<{ header?: { closeTime?: string } }> { - const response = await fetch(this.rpcUrl, { + const response = await this.egress.fetch(this.rpcUrl, { method: "POST", headers: { "Content-Type": "application/json" }, body: JSON.stringify({ @@ -51,11 +60,12 @@ export class SorobanService { method: "getLedgers", params: { startLedger: sequence, endLedger: sequence }, }), + purpose: EgressPurpose.RPC, }); - if (!response.ok) { - throw new Error(`getLedgers HTTP ${response.status} for ledger ${sequence}`); + if (response.statusCode < 200 || response.statusCode >= 300) { + throw new Error(`getLedgers HTTP ${response.statusCode} for ledger ${sequence}`); } - const body = (await response.json()) as { + const body = JSON.parse(response.body) as { result?: { ledgers?: Array<{ ledgerCloseTime?: string }> }; error?: { message?: string }; }; @@ -64,13 +74,6 @@ export class SorobanService { } const ledger = body.result?.ledgers?.[0]; return ledger ? { header: { closeTime: ledger.ledgerCloseTime } } : {}; - getLedger(sequence: number) { - // The stellar-sdk 12.x Server type no longer exposes `getLedger`; the call - // is preserved for the event-ingestion lag metric. Cast to keep compiling - // against the pinned SDK — the runtime API may need a follow-up migration. - return (this.server as unknown as { - getLedger(seq: number): Promise<{ header?: { closeTime?: string | number } }>; - }).getLedger(sequence); } getEvents(request: SorobanRpc.Server.GetEventsRequest) { diff --git a/src/treasury/treasury.service.ts b/src/treasury/treasury.service.ts index ea7d4c86..3ebac1a9 100644 --- a/src/treasury/treasury.service.ts +++ b/src/treasury/treasury.service.ts @@ -165,13 +165,10 @@ export class TreasuryService { const [code, issuer] = asset.split(":"); if (issuer) { const balance = account.balances.find( - (b) => - b.asset_type !== "native" && - "asset_code" in b && - "asset_issuer" in b && (b): b is typeof b & { asset_code: string; asset_issuer: string } => b.asset_type !== "native" && "asset_code" in b && + "asset_issuer" in b && b.asset_code === code && b.asset_issuer === issuer, ); diff --git a/test/chaos/runner.ts b/test/chaos/runner.ts index ec2be3a7..4f7d96d1 100644 --- a/test/chaos/runner.ts +++ b/test/chaos/runner.ts @@ -36,6 +36,7 @@ const results: ScenarioResult[] = []; // ── Toxiproxy helpers ──────────────────────────────────────────────────────── async function addToxic(proxy: string, toxic: ToxicConfig, name: string): Promise { + // eslint-disable-next-line no-restricted-syntax -- standalone chaos runner targets a local Toxiproxy endpoint const res = await fetch(`${TOXIPROXY}/proxies/${proxy}/toxics`, { method: "POST", headers: { "content-type": "application/json" }, @@ -54,6 +55,7 @@ async function addToxic(proxy: string, toxic: ToxicConfig, name: string): Promis } async function removeToxic(proxy: string, name: string): Promise { + // eslint-disable-next-line no-restricted-syntax -- standalone chaos runner targets a local Toxiproxy endpoint const res = await fetch(`${TOXIPROXY}/proxies/${proxy}/toxics/${name}`, { method: "DELETE", }); @@ -67,6 +69,7 @@ async function removeToxic(proxy: string, name: string): Promise { /** Returns the HTTP status of GET /health/ready. */ async function healthStatus(): Promise { try { + // eslint-disable-next-line no-restricted-syntax -- standalone chaos runner targets its configured test service const res = await fetch(`${BASE_URL}/health/ready`, { signal: AbortSignal.timeout(5_000) }); return res.status; } catch { @@ -81,6 +84,7 @@ async function healthStatus(): Promise { */ async function createTestIntent(): Promise { try { + // eslint-disable-next-line no-restricted-syntax -- standalone chaos runner targets its configured test service const res = await fetch(`${BASE_URL}/api/v1/intents`, { method: "POST", headers: {