diff --git a/.env.example b/.env.example index 4dd3c7d..e3e38c5 100644 --- a/.env.example +++ b/.env.example @@ -452,6 +452,10 @@ 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=[] + # Low-frequency scan that catches deadline jobs the queue did not run (issue #437). SAFETY_SWEEP_INTERVAL_MS=300000 diff --git a/CHANGELOG.md b/CHANGELOG.md index 7d0d367..4129c28 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). - Deadline jobs `expire-intent` and `fill-window-expired` replace the 30s sweeper poll. A safety sweep (`SAFETY_SWEEP_INTERVAL_MS`, default 5 min) catches lost jobs and increments `vortex_sweeper_safety_caught_total` (Closes #437). - `scripts/generate-client.ts` — generates a typed TypeScript API client from the live OpenAPI spec using `openapi-typescript` v7; output committed to `src/generated/` 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 e69de29..5c7b8f2 100644 --- a/jest.config.js +++ b/jest.config.js @@ -0,0 +1,63 @@ +/** @type {import('jest').Config} */ + +// The single definition of the coverage gate. It is enforced on the *merged* +// shard report by scripts/ci/coverage-merge.mjs, not by individual shard runs +// (issue #486), which pass --coverageThreshold '{}' because a shard only ever +// executes part of the suite. +module.exports = { + preset: "ts-jest", + testEnvironment: "node", + rootDir: "src", + testRegex: ".*\\.spec\\.ts$", + collectCoverageFrom: ["**/*.(t|j)s"], + coverageDirectory: "../coverage", + // text-summary keeps local runs readable, lcov feeds editors, and json is what + // coverage-merge.mjs consumes. + coverageReporters: ["text-summary", "lcov", "json"], + // Run both the main NestJS unit suite and the scripts suite under one command. + projects: [ + // ── Main NestJS unit suite ────────────────────────────────────────────── + { + displayName: "src", + preset: "ts-jest", + testEnvironment: "node", + rootDir: "src", + testRegex: ".*\\.spec\\.ts$", + // 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.) ───────────────────────────────── + // Tests live in src/scripts/ but import from scripts/ (outside src/). + // A dedicated tsconfig with broader rootDir handles the path. + { + displayName: "scripts", + testEnvironment: "node", + rootDir: ".", + testMatch: ["/src/scripts/**/*.spec.ts"], + transform: { + "^.+\\.tsx?$": [ + "ts-jest", + { + tsconfig: "./tsconfig.scripts.json", + }, + ], + }, + }, + ], + + // Coverage is collected from the project-level collectCoverageFrom above. + coverageDirectory: "coverage", + coverageThreshold: { + global: { + branches: 70, + functions: 70, + lines: 70, + statements: 70, + }, + }, +}; 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 3679138..e69de29 100644 --- a/prisma/schema.prisma +++ b/prisma/schema.prisma @@ -1,599 +0,0 @@ -// ─── Prisma Schema ──────────────────────────────────────────────────────────── -// Database: PostgreSQL (swap to sqlite for local dev / CI without a real DB) -// Run `npm run db:generate` after editing this file. -// ───────────────────────────────────────────────────────────────────────────── - -generator client { - provider = "prisma-client-js" -} - -datasource db { - provider = "postgresql" - url = env("DATABASE_URL") -} - -// ─── Enums ─────────────────────────────────────────────────────────────────── - -enum IntentState { - open - accepted - filled - cancelled - expired - slashed -} - -enum SupportedChain { - stellar - ethereum - base - polygon - arbitrum - optimism - avalanche -} - -/// Registry lifecycle. Delisted rows stay in the table so existing intents -/// that already copied the token metadata keep resolving. -enum TokenStatus { - active - paused - delisted - - @@map("token_status") -} - -// ─── Kill-switch scopes (issue #477) ────────────────────────────────────────── -// A switch is addressed by exactly one of four mutually-exclusive scopes. -// `global` has no chain/token; `chain` sets chain only; `token` sets chain + -// token; `operation` narrows further to a single operation name. -// Evaluation walks GLOBAL → CHAIN → TOKEN → OPERATION and the most specific -// matching switch wins, so an operation-level resume can never re-open a -// scope that a broader switch is still holding closed. - -enum KillSwitchScope { - global - chain - token - operation -} - -/// Operations a switch can gate. Constrained in the application layer (DTO + -/// normalisation) before it reaches the database, so the enum here is a -/// backstop rather than the only gate. -/// Changing this list requires a migration, since Prisma maps it to a real -/// Postgres enum type. -enum KillSwitchOperation { - create - accept - fill - slash - onchain -} - -// ─── Intent ────────────────────────────────────────────────────────────────── -// Represents a cross-chain swap request submitted by a user. -// The nested token objects (srcToken / dstToken) are stored as JSONB so the -// schema remains flexible while the intent's own scalar fields stay queryable. - -model Intent { - id String @id @default(uuid()) @map("id") - /// Externally-visible intent identifier (UUID). - intentId String @unique @map("intent_id") - /// Stellar / EVM address of the user that created the intent. - user String @map("user") - srcChain SupportedChain @map("src_chain") - /// ERC-20 / native token info on the source chain (stored as JSON). - srcToken Json @map("src_token") - /// Raw amount in the token's base unit (stored as string to preserve bigint precision). - srcAmount String @map("src_amount") - /// Stellar destination token info (stored as JSON). - dstToken Json @map("dst_token") - /// Minimum acceptable destination amount (base unit string). - minDstAmount String @map("min_dst_amount") - /// Best quote from solvers, populated after quoting. - quotedDstAmount String? @map("quoted_dst_amount") - /// Solver address that accepted / filled this intent. - solver String? @map("solver") - state IntentState @default(open) @map("state") - /// Unix epoch seconds. - createdAt Int @map("created_at") - /// Unix epoch seconds – intent expires after this. - deadline Int @map("deadline") - /// Unix epoch seconds – set when state becomes `filled`. - filledAt Int? @map("filled_at") - /// Actual amount received by the user (base unit string). - fillAmount String? @map("fill_amount") - /// Realized protocol fee charged on this fill (destination-token base units). - feeAmount String? @map("fee_amount") - /// On-chain transaction hash of the Stellar fill transaction. - txHash String? @map("tx_hash") - /// Unix epoch seconds – set when state becomes `slashed`. - slashedAt Int? @map("slashed_at") - slashReason String? @map("slash_reason") - /// Optimistic-concurrency version, incremented on every update (issue #405). - version Int @default(0) @map("version") - /// Caller-supplied idempotency key for POST /intents; unique across replicas (issue #404). - idempotencyKey String? @unique @map("idempotency_key") - /// Source-chain escrow deposit confirmed (issue #403); unverified intents are not fillable. - srcVerified Boolean @default(false) @map("src_verified") - /// EVM deposit transaction hash supplied at creation, if any. - srcTxHash String? @map("src_tx_hash") - /// Last verification result: status, block, received amount, detail. - srcVerification Json? @map("src_verification") - - @@index([user]) - @@index([state]) - @@index([solver]) - @@map("intents") -} - -// ─── Solver ────────────────────────────────────────────────────────────────── -// Registered solver nodes that fulfil cross-chain intents. - -model Solver { - id String @id @default(uuid()) @map("id") - /// Public key / address that uniquely identifies the solver. - address String @unique @map("address") - name String @map("name") - /// Bond amount locked in the registry contract (base-unit string). - bondAmount String @map("bond_amount") - fillsCompleted Int @default(0) @map("fills_completed") - fillsFailed Int @default(0) @map("fills_failed") - /// Cumulative filled volume (base-unit string). - totalVolume String @default("0") @map("total_volume") - /// Rolling average fill time in seconds. - avgFillTime Float @default(0) @map("avg_fill_time") - isActive Boolean @default(true) @map("is_active") - /// Unix epoch seconds. - registeredAt Int @map("registered_at") - /// Unix epoch seconds of the solver's most recent activity (registration, fill, or status change). - lastActiveAt Int @map("last_active_at") - /// Chains this solver supports (stored as JSON array of SupportedChain values). - supportedChains Json @map("supported_chains") - /// Token symbols this solver can handle. - supportedTokens Json @map("supported_tokens") - - // ── On-chain projection fields (issue #399) ─────────────────────────────── - /// Indicates the authoritative data source: "api" (REST registration) or - /// "chain" (projected from solver-registry contract events). - source String @default("api") @map("source") - /// Ledger sequence of the most recent on-chain event that updated this row. - /// NULL when source="api" (no on-chain event has been observed yet). - chainUpdatedLedger Int? @map("chain_updated_ledger") - - @@index([isActive]) - @@map("solvers") -} - -// ─── IntentAuditLog ────────────────────────────────────────────────────────── -// Append-only record of every state transition for an intent (issue #217 / #62). -// Queried by intentId to reconstruct the full history of a swap. - -model IntentAuditLog { - id BigInt @id @default(autoincrement()) @map("id") - /// FK to intents.intent_id (the user-visible UUID, not the surrogate PK). - intentId String @map("intent_id") - /// ISO-8601 / TIMESTAMPTZ of when the transition was recorded. - timestamp DateTime @default(now()) @map("timestamp") @db.Timestamptz - /// State the intent moved INTO (e.g. "cancelled", "expired", "slashed"). - toState String @map("to_state") - /// Actor who triggered the transition: a user address, solver address, or "system". - actor String @map("actor") - /// Human-readable explanation. - reason String @map("reason") - /// Optional extra data (fill amount, tx hash, deadline, …). - metadata Json? @map("metadata") - - @@index([intentId, timestamp(sort: Asc)], name: "audit_log_intent_idx") - @@map("intent_audit_log") -} - -// ─── ContractUpgrade ───────────────────────────────────────────────────────── -// Append-only history of detected WASM upgrades of the contracts the backend -// writes to (issue #402). Written by ContractVersionService when a poll sees -// a new hash or an upgrade event is ingested. - -model ContractUpgrade { - id BigInt @id @default(autoincrement()) @map("id") - /// "settlement" | "solverRegistry" - contractName String @map("contract_name") - contractId String @map("contract_id") - previousWasmHash String? @map("previous_wasm_hash") - wasmHash String? @map("wasm_hash") - /// ABI version the new hash maps to; null when unsupported. - abiVersion String? @map("abi_version") - /// "poll" | "event" - source String @map("source") - ledger Int? @map("ledger") - txHash String? @map("tx_hash") - detectedAt DateTime @default(now()) @map("detected_at") @db.Timestamptz - - @@index([contractId, detectedAt], name: "contract_upgrades_contract_idx") - @@map("contract_upgrades") -} - -// ─── Token ─────────────────────────────────────────────────────────────────── -// Static registry of tokens the protocol supports. -// Kept separate so it can be updated without migrations when the token list changes. - -model Token { - id String @id @default(uuid()) @map("id") - /// Contract address (EVM hex address or Stellar contract ID). - address String @map("address") - symbol String @map("symbol") - name String @map("name") - decimals Int @map("decimals") - chain SupportedChain @map("chain") - logoUri String? @map("logo_uri") - /// Latest known USD price (nullable – updated by a price-feed worker). - priceUsd Float? @map("price_usd") - /// Whether this is a destination-side Stellar token. - isStellar Boolean @default(false) @map("is_stellar") - /// active tokens are listed; paused stay listed but are not offered for new - /// intents; delisted are hidden from discovery and still readable by id. - status TokenStatus @default(active) @map("status") - /// evm | stellar-sac | stellar-classic. String so classic and SAC stay distinct - /// without a second database enum. - assetKind String @default("evm") @map("asset_kind") - - @@unique([address, chain]) - @@index([chain]) - @@index([symbol]) - @@map("tokens") -} - -// ─── KillSwitch (issue #477) ───────────────────────────────────────────────── -// Hierarchical emergency pause. One row per (scope, chain, token, operation) -// tuple; activating the row closes that scope for the gated operation. -// -// Because the unique key is the full scope tuple, "resume" is an idempotent -// state flip on a single row rather than a delete/insert race between replicas. -// -// `scope` is stored alongside the nullable chain/token/operation columns so the -// row is self-describing; the application maintains that invariant in one place -// (KillSwitchService.normaliseScope) rather than trusting callers. - -model KillSwitch { - id String @id @default(uuid()) @map("id") - scope KillSwitchScope @map("scope") - /// NULL for global scope. - chain String? @map("chain") - /// NULL for global and chain scopes. - token String? @map("token") - /// NULL unless scope = operation. - operation KillSwitchOperation? @map("operation") - - /// Canonical encoding of (scope, chain, token, operation), e.g. - /// "operation|stellar|USDC|fill". Postgres will not let Prisma build a - /// findUnique over nullable columns, so this non-null surrogate carries the - /// uniqueness guarantee while the typed columns above stay queryable. - scopeKey String @unique @map("scope_key") - - /// true = writes are blocked for this scope, false = explicitly resumed. - active Boolean @map("active") - /// Machine-readable reason surfaced to clients as `reason` in the 503 body. - reasonCode String @map("reason_code") - /// Free-text operator explanation (incident ticket, etc). - reason String @map("reason") - /// Operator identity that activated or resumed the switch. - activatedBy String @map("activated_by") - - /// Unix epoch ms of the last state change. Used as the cheap change-detection - /// key for the polling fallback (monotonic: re-pausing bumps it). - updatedAt Int @map("updated_at") - - createdAt Int @map("created_at") - /// Unix epoch ms of the last time this switch transitioned active -> inactive. - /// Unset while never resumed, so a brand-new row cannot inherit a stale - /// cooldown from a previous pause of the same scope. - lastResumedAt Int? @map("last_resumed_at") - - // Resume requires two distinct approvals; see KillSwitchApproval. - approvals KillSwitchApproval[] - /// Number of distinct approvals still required before the switch may resume. - approvalsRequired Int @default(2) @map("approvals_required") - - @@index([active]) - @@index([scope, chain, token, operation]) - /// Backs the polling fallback's "max(updated_at) since last snapshot?" probe. - @@index([updatedAt]) - @@map("kill_switches") -} -// ─── KillSwitchApproval (issue #477) ────────────────────────────────────────── -// One row per operator approval to resume a switch. The resume guard is a -// count of DISTINCT approver over this table, so a single operator cannot -// approve twice, and two different operators can never be one person unless the -// caller lies about identity (the API ties identity to the operator token). - -model KillSwitchApproval { - id String @id @default(uuid()) @map("id") - killSwitchId String @map("kill_switch_id") - killSwitch KillSwitch @relation(fields: [killSwitchId], references: [id], onDelete: Cascade) - /// Operator identity that granted this approval. Unique per switch so one - /// operator cannot satisfy the two-approval rule alone. - approver String @map("approver") - /// Unix epoch ms. - approvedAt Int @map("approved_at") - /// Optional note explaining the approval. - note String? @map("note") - - @@unique([killSwitchId, approver]) - @@index([killSwitchId]) - @@map("kill_switch_approvals") -} - -// ─── TreasurySnapshot ──────────────────────────────────────────────────────── -// Daily snapshots of expected vs actual treasury balances per asset. -// Used for reconciliation and historical reporting. - -model TreasurySnapshot { - id BigInt @id @default(autoincrement()) @map("id") - /// Date of the snapshot (YYYY-MM-DD). - snapshotDate String @map("snapshot_date") - /// Asset identifier (contract address or asset code). - asset String @map("asset") - /// Expected balance from fee ledger + slashes - refunds (base-unit string). - expectedBalance String @map("expected_balance") - /// Actual on-chain balance (base-unit string). - actualBalance String @map("actual_balance") - /// Difference (actualBalance - expectedBalance, base-unit string). - discrepancy String @map("discrepancy") - /// Tolerance threshold used for this check (base-unit string). - toleranceThreshold String @map("tolerance_threshold") - /// Whether abs(discrepancy) > toleranceThreshold. - hasUnexplainedDiscrepancy Boolean @map("has_unexplained_discrepancy") - /// Human-readable explanation of differences (in-flight settlements, refunds, etc). - explanation String? @map("explanation") - /// ISO-8601 / TIMESTAMPTZ of when the snapshot was taken. - createdAt DateTime @default(now()) @map("created_at") @db.Timestamptz - /// Breakdown of expected balance sources (fees, slashes, refunds). - breakdown Json? @map("breakdown") - - @@unique([snapshotDate, asset], name: "snapshot_date_asset_unique") - @@index([snapshotDate]) - @@index([hasUnexplainedDiscrepancy]) - @@map("treasury_snapshots") -} - -// ─── FeeLedger ─────────────────────────────────────────────────────────────── -// Append-only ledger of all fee accruals from filled intents. - -model FeeLedger { - id BigInt @id @default(autoincrement()) @map("id") - /// FK to intents.intent_id. - intentId String @map("intent_id") - /// Asset identifier (contract address or asset code). - asset String @map("asset") - /// Fee amount in base units (string). - amount String @map("amount") - /// ISO-8601 / TIMESTAMPTZ of when the fee was accrued. - accrualAt DateTime @default(now()) @map("accrual_at") @db.Timestamptz - /// Optional transaction hash. - txHash String? @map("tx_hash") - - @@index([intentId]) - @@index([asset, accrualAt]) - @@map("fee_ledger") -} - -// ─── SlashLedger ───────────────────────────────────────────────────────────── -// Append-only ledger of all slash proceeds from solver penalties. - -model SlashLedger { - id BigInt @id @default(autoincrement()) @map("id") - /// Solver address that was slashed. - solverAddress String @map("solver_address") - /// Asset identifier (contract address or asset code). - asset String @map("asset") - /// Slashed amount in base units (string). - amount String @map("amount") - /// ISO-8601 / TIMESTAMPTZ of when the slash occurred. - slashedAt DateTime @default(now()) @map("slashed_at") @db.Timestamptz - /// Reason for the slash. - reason String @map("reason") - /// Optional transaction hash. - txHash String? @map("tx_hash") - - @@index([solverAddress]) - @@index([asset, slashedAt]) - @@map("slash_ledger") -} - -// ─── RefundLedger ──────────────────────────────────────────────────────────── -// Append-only ledger of all refunds issued to users. - -model RefundLedger { - id BigInt @id @default(autoincrement()) @map("id") - /// FK to intents.intent_id. - intentId String @map("intent_id") - /// User address that received the refund. - userAddress String @map("user_address") - /// Asset identifier (contract address or asset code). - asset String @map("asset") - /// Refund amount in base units (string). - amount String @map("amount") - /// ISO-8601 / TIMESTAMPTZ of when the refund was issued. - issuedAt DateTime @default(now()) @map("issued_at") @db.Timestamptz - /// Reason for the refund. - reason String @map("reason") - /// Optional transaction hash. - txHash String? @map("tx_hash") - - @@index([intentId]) - @@index([userAddress]) - @@index([asset, issuedAt]) - @@map("refund_ledger") -} - -// ─── AdminAuditLog ─────────────────────────────────────────────────────────── -// Append-only record of privileged operator and governance actions: feature -// flag changes (issue #495), kill-switch toggles and guardian overrides (#507). - -model AdminAuditLog { - id BigInt @id @default(autoincrement()) @map("id") - /// Admin principal id, or "guardian" / "system" for automated actions. - actor String @map("actor") - /// Dotted action name, e.g. "flag.update", "guardian.override". - action String @map("action") - /// Target of the action, e.g. "flag:onchain-dry-run". - target String @map("target") - before Json? @map("before") - after Json? @map("after") - reason String? @map("reason") - createdAt DateTime @default(now()) @map("created_at") @db.Timestamptz - - @@index([target, createdAt(sort: Desc)]) - @@map("admin_audit_log") -} - -// ─── FeatureFlag ───────────────────────────────────────────────────────────── -// Runtime feature flags (issue #495). `rules` is an ordered JSON array of -// { value, percentage?, solvers?, chains? }; the first matching rule wins. - -model FeatureFlag { - key String @id @map("key") - defaultValue Boolean @map("default_value") - rules Json @map("rules") - version Int @default(1) @map("version") - updatedBy String @map("updated_by") - updatedAt DateTime @updatedAt @map("updated_at") @db.Timestamptz - - @@map("feature_flags") -} - -// Pending flag changes that need a second approver (e.g. dry-run off in production). -model FlagChangeRequest { - id String @id @default(uuid()) @map("id") - flagKey String @map("flag_key") - /// Proposed { defaultValue, rules }. - proposed Json @map("proposed") - proposedBy String @map("proposed_by") - /// Admin ids that approved, proposer first. - approvals String[] @map("approvals") - /// pending | applied - status String @default("pending") @map("status") - reason String? @map("reason") - createdAt DateTime @default(now()) @map("created_at") @db.Timestamptz - appliedAt DateTime? @map("applied_at") @db.Timestamptz - - @@index([flagKey, status]) - @@map("flag_change_requests") -} - -// ─── GuardianAction ────────────────────────────────────────────────────────── -// Emergency actions ingested from the on-chain guardian contract (issue #507). -// One row per activating event; cleared by the matching guardian event or a -// superadmin override. - -model GuardianAction { - /// Soroban event id of the activating event. - id String @id @map("id") - /// pause | freeze | blacklist - kind String @map("kind") - /// Frozen parameter key or blacklisted solver address; "" for pause. - target String @map("target") - active Boolean @map("active") - txHash String @map("tx_hash") - ledger Int @map("ledger") - activatedAt DateTime @map("activated_at") @db.Timestamptz - clearedAt DateTime? @map("cleared_at") @db.Timestamptz - clearedTxHash String? @map("cleared_tx_hash") - overriddenBy String? @map("overridden_by") - - @@index([active]) - @@map("guardian_actions") -} - -// ─── OnchainOutbox ─────────────────────────────────────────────────────────── -// Transactional outbox for on-chain writes (issue #396). A row is inserted in -// the same database transaction as the intent mutation it mirrors, and -// OutboxRelayService submits it to Soroban afterwards — so the DB can never say -// "accepted" while the corresponding transaction silently never went out. -// -// `id` is a monotonically increasing sequence and doubles as the per-intent -// ordering key: a row is only claimable once every earlier row for the same -// intent has reached a terminal success state. - -enum OutboxStatus { - pending - processing - submitted - confirmed - simulated - dead -} - -model OnchainOutbox { - id BigInt @id @default(autoincrement()) @map("id") - /// intents.intent_id this operation belongs to (ordering partition key). - intentId String @map("intent_id") - /// Contract operation, e.g. "create_intent", "accept_intent". - operation String @map("operation") - /// Operation arguments (JSON-safe, bigint amounts as strings). - payload Json @map("payload") - status OutboxStatus @default(pending) @map("status") - /// Number of times this row has been claimed by a relay worker. - attempts Int @default(0) @map("attempts") - nextAttemptAt DateTime @default(now()) @map("next_attempt_at") @db.Timestamptz - /// Lease expiry while status = processing; a crashed worker's row is reclaimed after this. - lockedUntil DateTime? @map("locked_until") @db.Timestamptz - /// Hash of the signed envelope, persisted BEFORE submission (crash idempotency). - envelopeHash String? @map("envelope_hash") - /// Hash of the transaction that was submitted. - txHash String? @map("tx_hash") - lastError String? @map("last_error") - createdAt DateTime @default(now()) @map("created_at") @db.Timestamptz - updatedAt DateTime @default(now()) @updatedAt @map("updated_at") @db.Timestamptz - - @@index([status, nextAttemptAt], name: "onchain_outbox_claim_idx") - @@index([intentId, id], name: "onchain_outbox_intent_order_idx") - @@map("onchain_outbox") -} - -// ─── PendingSlash ──────────────────────────────────────────────────────────── -// Durable saga state for a solver slash (issue #397): -// detected → challenge_window → submitted → confirmed | cancelled -// The unique constraint on intent_id enforces exactly-once slashing per intent. - -enum PendingSlashState { - detected - challenge_window - submitted - confirmed - cancelled -} - -model PendingSlash { - id String @id @default(uuid()) @map("id") - intentId String @unique @map("intent_id") - solverAddress String @map("solver_address") - reason String @map("reason") - state PendingSlashState @default(detected) @map("state") - /// The intent's fill deadline (unix seconds) that the solver missed. - fillDeadline Int @map("fill_deadline") - detectedAt DateTime @map("detected_at") @db.Timestamptz - /// Slash may not be submitted before this instant. - challengeEndsAt DateTime @map("challenge_ends_at") @db.Timestamptz - attempts Int @default(0) @map("attempts") - nextAttemptAt DateTime @default(now()) @map("next_attempt_at") @db.Timestamptz - /// Lease held by the pipeline worker that is currently verifying/submitting. - lockedUntil DateTime? @map("locked_until") @db.Timestamptz - txHash String? @map("tx_hash") - /// true when the registry client simulated but did not broadcast (dry-run / gated submit). - simulated Boolean @default(false) @map("simulated") - submittedAt DateTime? @map("submitted_at") @db.Timestamptz - confirmedAt DateTime? @map("confirmed_at") @db.Timestamptz - cancelledAt DateTime? @map("cancelled_at") @db.Timestamptz - cancelReason String? @map("cancel_reason") - cancelledBy String? @map("cancelled_by") - /// Fill transaction that cancelled the slash, when cancelled by a fill proof. - fillTxHash String? @map("fill_tx_hash") - lastError String? @map("last_error") - createdAt DateTime @default(now()) @map("created_at") @db.Timestamptz - updatedAt DateTime @default(now()) @updatedAt @map("updated_at") @db.Timestamptz - - @@index([state, challengeEndsAt], name: "pending_slashes_due_idx") - @@index([solverAddress], name: "pending_slashes_solver_idx") - @@map("pending_slashes") -} 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/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/stats/stats.service.ts b/src/stats/stats.service.ts index bfc93b2..e69de29 100644 --- a/src/stats/stats.service.ts +++ b/src/stats/stats.service.ts @@ -1,215 +0,0 @@ -import { Injectable, Optional } from "@nestjs/common"; -import { ConfigService } from "@nestjs/config"; -import { createHmac } from "crypto"; -import { AppConfig } from "../config/configuration"; -import { isCanaryIntent } from "../common/canary"; -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 { AbuseScoreService } from "../abuse/abuse-score.service"; -import { CHALLENGE_THRESHOLD } from "../abuse/abuse.types"; - -@Injectable() -export class StatsService { - constructor( - private readonly intentsService: IntentsService, - private readonly solversService: SolversService, - private readonly intentsGateway: IntentsGateway, - @Optional() private readonly abuseScorer?: AbuseScoreService, - @Optional() config?: ConfigService, - ) { - this.canary = new Set(config?.get("canaryAddresses", { infer: true }) ?? []); - } - - /** Canary addresses (issue #496) — their intents and solvers never count toward public stats. */ - private readonly canary: ReadonlySet; - - private async publicIntents() { - return (await this.intentsService.getAll()).filter((i) => !isCanaryIntent(i, this.canary)); - } - - /** - * Returns intents for stat computation, optionally excluding flagged (high-abuse-score) actors. - * - * When `includeFlagged` is false (the default), intents from users whose cached - * abuse score is above the challenge threshold are excluded so spam does not - * inflate governance and incentive metrics. - * - * Operators can toggle `includeFlagged=true` on any stats endpoint for a - * full-transparency view that shows the raw counts before filtering. - */ - private async publicIntentsFiltered(includeFlagged = false) { - const intents = await this.publicIntents(); - if (includeFlagged || !this.abuseScorer) return intents; - - // Fetch cached scores for all unique users in parallel (best-effort: missing - // score = unknown = keep). - const uniqueUsers = [...new Set(intents.map((i) => i.user))]; - const scores = await Promise.all( - uniqueUsers.map(async (u) => [u, await this.abuseScorer!.getCachedScore(u)] as const), - ); - const flaggedUsers = new Set( - scores.filter(([, score]) => score !== null && score >= CHALLENGE_THRESHOLD).map(([u]) => u), - ); - - return intents.filter((i) => !flaggedUsers.has(i.user)); - } - - private canonicalJson(value: unknown): string { - const seen = new WeakSet(); - - const normalize = (input: unknown): unknown => { - if (Array.isArray(input)) { - return input.map((item) => normalize(item)); - } - if (input && typeof input === "object") { - if (seen.has(input)) { - return "[Circular]"; - } - seen.add(input); - return Object.keys(input as Record) - .sort() - .reduce>((acc, key) => { - acc[key] = normalize((input as Record)[key]); - return acc; - }, {}); - } - return input; - }; - - return JSON.stringify(normalize(value)); - } - - private getPublicDatasetUrl() { - return process.env.PUBLIC_STATS_DATASET_URL ?? "https://example.invalid/public-stats/latest.json"; - } - - private getPublicSigningKey() { - return process.env.PUBLIC_STATS_SIGNING_KEY ?? "public-stats-dev-key"; - } - - private getProvenance(payload: Record) { - const safePayload = { ...payload }; - delete safePayload.provenance; - const generatedAt = new Date().toISOString(); - const signature = createHmac("sha256", this.getPublicSigningKey()) - .update(this.canonicalJson({ ...safePayload, generatedAt })) - .digest("hex"); - - return { - watermarkLedger: 0, - generatedAt, - queryVersion: "public-stats/v1", - datasetUrl: this.getPublicDatasetUrl(), - signature, - }; - } - - async getProtocolStats(includeFlagged = false) { - const intents = await this.publicIntentsFiltered(includeFlagged); - const solvers = (await this.solversService.getAll()).filter((s) => !this.canary.has(s.address)); - - const open = intents.filter((i) => i.state === "open").length; - const filled = intents.filter((i) => i.state === "filled"); - const totalVolume = filled.reduce((sum, i) => sum + BigInt(i.fillAmount ?? "0"), 0n); - - const fillTimes = filled - .filter((i) => i.filledAt != null) - .map((i) => i.filledAt! - i.createdAt); - const avgFillTime = fillTimes.length - ? fillTimes.reduce((a, b) => a + b, 0) / fillTimes.length - : 0; - - return { - totalIntents: intents.length, - openIntents: open, - totalVolume: totalVolume.toString(), - uniqueUsers: new Set(intents.map((i) => i.user)).size, - activeSolvers: solvers.filter((s) => s.isActive).length, - avgFillTime: Math.round(avgFillTime), - fillRate: intents.length ? filled.length / intents.length : 0, - flaggedActivityExcluded: !includeFlagged, - }; - } - - async getPublicStats(includeFlagged = false) { - const stats = await this.getProtocolStats(includeFlagged); - return { - ...stats, - provenance: this.getProvenance(stats), - }; - } - - getPublicStatsHistory() { - return [] as Array>; - } - - async getTreasuryStats(includeFlagged = false) { - const intents = await this.publicIntentsFiltered(includeFlagged); - const now = Math.floor(Date.now() / 1000); - const last24hCutoff = now - 86_400; - - const allTime = intents - .filter((intent) => typeof intent.feeAmount === "string" && intent.feeAmount.length > 0) - .reduce((sum, intent) => sum + BigInt(intent.feeAmount ?? "0"), 0n); - - const last24h = intents - .filter( - (intent) => - typeof intent.feeAmount === "string" && - intent.feeAmount.length > 0 && - typeof intent.filledAt === "number" && - intent.filledAt >= last24hCutoff, - ) - .reduce((sum, intent) => sum + BigInt(intent.feeAmount ?? "0"), 0n); - - const byChain = new Map(); - - for (const intent of intents) { - if (typeof intent.feeAmount !== "string" || intent.feeAmount.length === 0) continue; - const fee = BigInt(intent.feeAmount ?? "0"); - const entry = byChain.get(intent.srcChain) ?? { - totalFees: 0n, - last24hFees: 0n, - filledCount: 0, - }; - - entry.totalFees += fee; - entry.filledCount += 1; - if (typeof intent.filledAt === "number" && intent.filledAt >= last24hCutoff) { - entry.last24hFees += fee; - } - byChain.set(intent.srcChain, entry); - } - - return { - allTime: { - totalFees: allTime.toString(), - filledIntents: intents.filter((intent) => typeof intent.feeAmount === "string" && intent.feeAmount.length > 0).length, - }, - last24h: { - totalFees: last24h.toString(), - filledIntents: intents.filter( - (intent) => - typeof intent.feeAmount === "string" && - intent.feeAmount.length > 0 && - typeof intent.filledAt === "number" && - intent.filledAt >= last24hCutoff, - ).length, - }, - byChain: Array.from(byChain.entries()).map(([srcChain, stats]) => ({ - srcChain, - totalFees: stats.totalFees.toString(), - last24hFees: stats.last24hFees.toString(), - filledIntents: stats.filledCount, - })), - }; - } - - getWsStats() { - return { - subscriberCount: this.intentsGateway.getSubscriberCount(), - }; - } -}