From 47727e3898b3af070a7eb8ddf762c9610f516d22 Mon Sep 17 00:00:00 2001 From: Goodnessukaigwe Date: Tue, 29 Sep 2026 10:09:19 +0100 Subject: [PATCH] feat(intents): expire intents from deadline jobs instead of a 30s poll Arm expire-intent and fill-window-expired at create, accept, and deadline extension, and keep a slower safety sweep for jobs the queue misses. Co-authored-by: Cursor --- .env.example | 11 ++ CHANGELOG.md | 1 + docs/adr/0005-deadline-jobs.md | 25 +++ docs/runbooks/on-call.md | 4 +- jest.config.js | 3 + package-lock.json | 160 +++++++++++++++++- package.json | 3 + src/app.module.ts | 23 --- src/common/stellar-signature.ts | 3 + src/config/configuration.ts | 24 +++ src/config/env.validation.ts | 14 +- src/governance/governance.module.ts | 4 +- src/health/health-indicator.registry.ts | 2 +- src/intents/intents-deadline.jobs.spec.ts | 152 +++++++++++++++++ src/intents/intents-deadline.jobs.ts | 83 +++++++++ .../intents-sweeper.manual-trigger.spec.ts | 2 +- src/intents/intents-sweeper.service.ts | 119 ++++++++----- src/intents/intents.gateway.spec.ts | 23 ++- src/intents/intents.gateway.ts | 8 +- src/intents/intents.module.ts | 10 +- src/intents/intents.service.shadow.spec.ts | 11 ++ src/intents/intents.service.spec.ts | 17 +- src/intents/intents.service.ts | 24 ++- src/intents/solver-intent-matcher.ts | 19 ++- src/intents/ws/connection-state.ts | 16 +- src/metrics/metrics.service.ts | 13 ++ 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 | 20 +++ 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/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 +- 42 files changed, 798 insertions(+), 336 deletions(-) create mode 100644 docs/adr/0005-deadline-jobs.md create mode 100644 src/intents/intents-deadline.jobs.spec.ts create mode 100644 src/intents/intents-deadline.jobs.ts create mode 100644 test/__mocks__/nestjs-schedule.ts diff --git a/.env.example b/.env.example index 6ede42f..10af6c5 100644 --- a/.env.example +++ b/.env.example @@ -311,3 +311,14 @@ 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= +# Low-frequency scan that catches deadline jobs the queue did not run (issue #437). +SAFETY_SWEEP_INTERVAL_MS=300000 +# 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..7d0d367 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 +- 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/` (Closes #134) diff --git a/docs/adr/0005-deadline-jobs.md b/docs/adr/0005-deadline-jobs.md new file mode 100644 index 0000000..dc8a83c --- /dev/null +++ b/docs/adr/0005-deadline-jobs.md @@ -0,0 +1,25 @@ +# ADR 0005: Deadline jobs for intent expiry + +- **Status**: Accepted +- **Date**: 2026-09-29 +- **Technical Story**: #437 — replace the polling sweeper with deadline-scheduled jobs + +## Context + +`IntentsSweeperService` scanned every open and accepted intent every 30 seconds. Expiry latency was bounded by that interval, and the scan grew with intent volume. + +The job queue from ADR 0002 is already in the process. It has delayed enqueue and idempotency keys. It does not have a cancel or replace API. + +## Decision + +Arm `expire-intent` when an intent is created and `fill-window-expired` when it is accepted or its fill deadline is extended. The delay is the time remaining until the stored deadline. The idempotency key is `job:intentId:deadline`. + +A moved deadline enqueues a new job. The previous job still runs. The handler loads the intent and returns without writing when the state is terminal or the stored deadline is not the one in the payload. + +The leader-elected sweep stays, at `SAFETY_SWEEP_INTERVAL_MS` (default 5 minutes), for jobs lost to a crash or a missed timer. Each intent it expires or slashes increments `vortex_sweeper_safety_caught_total`. That counter should stay near zero. `triggerManualSweep` is unchanged. + +## Consequences + +- Expiry no longer waits for the scan interval. A job fires at the deadline. +- Stale jobs are ignored rather than cancelled. +- Operators tune the safety net with `SAFETY_SWEEP_INTERVAL_MS`. The queue driver remains `JOBS_DRIVER` (`memory` or `bullmq`); this change does not add a second queue. diff --git a/docs/runbooks/on-call.md b/docs/runbooks/on-call.md index 276f33b..2b938d2 100644 --- a/docs/runbooks/on-call.md +++ b/docs/runbooks/on-call.md @@ -187,7 +187,7 @@ A sweep that has been delayed or killed will simply be absent. ### How the sweeper works `IntentsSweeperService.sweep()` is triggered by a `setInterval` every -`SWEEP_INTERVAL_MS` (30 000 ms, hardcoded). It: +`SAFETY_SWEEP_INTERVAL_MS` (default 300 000 ms). Deadline jobs (`expire-intent`, `fill-window-expired`) are the primary path. The safety sweep: 1. Calls `IntentsService.getByState("open")` — iterates the in-memory store. 2. Compares each intent's `deadline` (Unix timestamp) against `Date.now()`. @@ -523,7 +523,7 @@ handlers never wait for Redis. | `STELLAR_NETWORK` | `testnet` | Network passphrase selection | | `PORT` | `4000` | HTTP + WS listen port | | `NODE_ENV` | `development` | Log verbosity (set to `production` in prod) | -| `SWEEP_INTERVAL_MS` | `30000` (hardcoded) | How often the sweeper runs; change requires code deploy | +| `SAFETY_SWEEP_INTERVAL_MS` | `300000` | How often the safety sweep scans for deadline jobs the queue missed | | `KILLSWITCH_OPERATOR_TOKEN` | empty (control plane disabled) | Secret for `/api/v1/ops/killswitch`; **required in production** | | `KILLSWITCH_REDIS_URL` | `REDIS_URL` when `WS_BACKPLANE=redis` | Cross-replica pause propagation; empty = poll only | | `KILLSWITCH_POLL_MS` | `2000` | DB change-probe interval backing up Redis; caps propagation delay | 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/src/app.module.ts b/src/app.module.ts index b93ab1d..882100a 100644 --- a/src/app.module.ts +++ b/src/app.module.ts @@ -11,20 +11,15 @@ 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"; - @Module({ imports: [ // Issue #44 — global rate limit: 100 requests per 60 s per IP @@ -35,7 +30,6 @@ import { DatasetsModule } from "./datasets/datasets.module"; limit: 100, }, ]), - // Enable scheduled tasks (cron jobs) ScheduleModule.forRoot(), ConfigModule, PrismaModule, @@ -43,23 +37,8 @@ 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, @@ -71,13 +50,11 @@ import { DatasetsModule } from "./datasets/datasets.module"; 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..6a13930 100644 --- a/src/config/configuration.ts +++ b/src/config/configuration.ts @@ -281,6 +281,19 @@ export interface AppConfig { /** Soroban RPC endpoints probed for quorum (majority must be healthy). */ rpcHealthUrls: string[]; }; + /** Low-frequency sweeper that catches deadline jobs the queue did not run. */ + safetySweepIntervalMs: number; + /** 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 +408,17 @@ export default (): AppConfig => ({ .map((u) => u.trim()) .filter(Boolean), }, + safetySweepIntervalMs: parseInt(process.env.SAFETY_SWEEP_INTERVAL_MS ?? "300000", 10), + 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..2adb21f 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,16 @@ 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(""), + SAFETY_SWEEP_INTERVAL_MS: Joi.number().integer().min(1000).default(300_000), + 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/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/intents-deadline.jobs.spec.ts b/src/intents/intents-deadline.jobs.spec.ts new file mode 100644 index 0000000..1320761 --- /dev/null +++ b/src/intents/intents-deadline.jobs.spec.ts @@ -0,0 +1,152 @@ +import { ConfigService } from "@nestjs/config"; +import { IntentsService } from "./intents.service"; +import { IntentsSweeperService } from "./intents-sweeper.service"; +import { IntentDeadlineScheduler, delayUntil } from "./intents-deadline.jobs"; +import { IntentsGateway } from "./intents.gateway"; +import { JobsService } from "../jobs/jobs.service"; +import { AppConfig } from "../config/configuration"; +import { InMemoryIntentsRepository } from "./intents.repository"; +import { PrismaService } from "../prisma/prisma.service"; +import { ProtocolParamsService } from "../governance/params.service"; +import { StellarTxService } from "../soroban/stellar-tx.service"; +import { MetricsService } from "../metrics/metrics.service"; +import { KillSwitchService } from "../killswitch/killswitch.service"; +import { SolversService } from "../solvers/solvers.service"; +import { SolverRegistryService } from "../soroban/solver-registry.service"; +import { LeaderElectionService } from "../common/leader-election"; + +function jobsConfig(): ConfigService { + return { + get: (key: string) => { + if (key === "jobs") return { driver: "memory", shutdownTimeoutMs: 1000 }; + if (key === "processRole") return "all"; + if (key === "safetySweepIntervalMs") return 300_000; + return undefined; + }, + } as unknown as ConfigService; +} + +function buildIntents(scheduler: IntentDeadlineScheduler): IntentsService { + const repo = new InMemoryIntentsRepository(); + (repo as unknown as { store: Map }).store.clear(); + const protocolParams = { + snapshotForChain: jest.fn().mockReturnValue({ + version: 0, + feeBps: 30, + deadlineSeconds: 1800, + fillWindowSeconds: 600, + capturedAt: new Date().toISOString(), + }), + } as unknown as ProtocolParamsService; + return new IntentsService( + repo, + { get: jest.fn().mockReturnValue(false) } as unknown as ConfigService, + {} as StellarTxService, + { intentAuditLog: { create: jest.fn().mockResolvedValue({}), findMany: jest.fn().mockResolvedValue([]) } } as unknown as PrismaService, + protocolParams, + undefined, + undefined, + undefined, + scheduler, + ); +} + +describe("deadline jobs", () => { + let jobs: JobsService; + let sweeper: IntentsSweeperService; + let intents: IntentsService; + const recordSafetyCatch = jest.fn(); + + beforeEach(() => { + jobs = new JobsService(jobsConfig()); + const scheduler = new IntentDeadlineScheduler(jobs); + intents = buildIntents(scheduler); + const metrics = { recordSweep: jest.fn(), recordSafetyCatch } as unknown as MetricsService; + sweeper = new IntentsSweeperService( + intents, + { broadcast: jest.fn().mockResolvedValue(undefined) } as unknown as IntentsGateway, + {} as SolversService, + { slashSolver: jest.fn().mockResolvedValue({ detail: "no-op" }) } as unknown as SolverRegistryService, + metrics, + { evaluateTarget: jest.fn().mockReturnValue({ paused: false, matched: null, matchedChain: [] }) } as unknown as KillSwitchService, + { registerWorker: jest.fn() } as unknown as LeaderElectionService, + jobsConfig(), + jobs, + ); + sweeper.onModuleInit(); + recordSafetyCatch.mockClear(); + }); + + afterEach(async () => { + await jobs.onApplicationShutdown(); + }); + + it("expires an open intent within 2s of its deadline", async () => { + const deadline = Math.floor(Date.now() / 1000) + 1; + const intent = await intents.create({ + user: "GTEST", + srcChain: "stellar", + srcToken: { address: "native", symbol: "XLM", name: "Stellar Lumens", decimals: 7, chain: "stellar" }, + srcAmount: "1", + dstToken: { contract: "CTEST", symbol: "USDC", decimals: 7 }, + minDstAmount: "1", + deadline, + }); + + const started = Date.now(); + let state = (await intents.get(intent.intentId))?.state; + while (state !== "expired" && Date.now() - started < 2500) { + await new Promise((resolve) => setTimeout(resolve, 20)); + state = (await intents.get(intent.intentId))?.state; + } + expect(state).toBe("expired"); + expect(Date.now() - deadline * 1000).toBeLessThan(2000); + }); + + it("ignores a stale expire job when the deadline was moved", async () => { + const past = Math.floor(Date.now() / 1000) - 5; + const intent = await intents.create({ + user: "GTEST", + srcChain: "stellar", + srcToken: { address: "native", symbol: "XLM", name: "Stellar Lumens", decimals: 7, chain: "stellar" }, + srcAmount: "1", + dstToken: { contract: "CTEST", symbol: "USDC", decimals: 7 }, + minDstAmount: "1", + deadline: past + 10_000, + }); + await intents.update(intent.intentId, { deadline: past }); + + await sweeper.handleExpireJob({ intentId: intent.intentId, deadline: past + 10_000 }); + expect((await intents.get(intent.intentId))?.state).toBe("open"); + + await sweeper.handleExpireJob({ intentId: intent.intentId, deadline: past }); + expect((await intents.get(intent.intentId))?.state).toBe("expired"); + await sweeper.handleExpireJob({ intentId: intent.intentId, deadline: past }); + expect((await intents.get(intent.intentId))?.state).toBe("expired"); + }); + + it("counts intents the safety sweep still has to settle", async () => { + await sweeper.sweep({ safety: true }); + expect(recordSafetyCatch).toHaveBeenCalledWith(0); + + const past = Math.floor(Date.now() / 1000) - 5; + const late = await intents.create({ + user: "GTEST2", + srcChain: "stellar", + srcToken: { address: "native", symbol: "XLM", name: "Stellar Lumens", decimals: 7, chain: "stellar" }, + srcAmount: "1", + dstToken: { contract: "CTEST", symbol: "USDC", decimals: 7 }, + minDstAmount: "1", + deadline: Math.floor(Date.now() / 1000) + 3600, + }); + await intents.update(late.intentId, { deadline: past }); + recordSafetyCatch.mockClear(); + await sweeper.sweep({ safety: true }); + expect(recordSafetyCatch).toHaveBeenCalledWith(1); + }); + + it("delayUntil is zero once the deadline has passed", () => { + expect(delayUntil(1, 5_000)).toBe(0); + expect(delayUntil(10, 1_000)).toBe(9_000); + }); +}); diff --git a/src/intents/intents-deadline.jobs.ts b/src/intents/intents-deadline.jobs.ts new file mode 100644 index 0000000..dd3dfdc --- /dev/null +++ b/src/intents/intents-deadline.jobs.ts @@ -0,0 +1,83 @@ +import { Injectable, Logger, Optional } from "@nestjs/common"; +import { JobsService } from "../jobs/jobs.service"; +import { defineJob } from "../jobs/jobs.types"; +import { Intent } from "./intents.types"; + +export const DEADLINE_QUEUE = "intents-deadlines"; + +/** Payload for a deadline job. `deadline` is the unix-second instant this job was armed for. */ +export interface DeadlineJobData { + intentId: string; + deadline: number; +} + +/** Fires when an open intent's user deadline is reached. */ +export const EXPIRE_INTENT_JOB = defineJob(DEADLINE_QUEUE, "expire-intent", { + attempts: 3, + backoffMs: 250, +}); + +/** Fires when an accepted intent's fill window is reached. */ +export const FILL_WINDOW_EXPIRED_JOB = defineJob(DEADLINE_QUEUE, "fill-window-expired", { + attempts: 3, + backoffMs: 250, +}); + +/** Milliseconds until `deadlineSec`, never negative. Exported for tests. */ +export function delayUntil(deadlineSec: number, nowMs = Date.now()): number { + return Math.max(0, deadlineSec * 1000 - nowMs); +} + +/** + * Enqueues deadline jobs. There is no cancel API on the queue, so a moved + * deadline is a new idempotency key (`kind:intentId:deadline`). The previous + * job still runs and the handler ignores it when the stored deadline differs. + */ +@Injectable() +export class IntentDeadlineScheduler { + private readonly logger = new Logger(IntentDeadlineScheduler.name); + + constructor(@Optional() private readonly jobs?: JobsService) {} + + /** Arm `expire-intent` for an open intent. */ + scheduleExpire(intent: Pick): void { + this.enqueue(EXPIRE_INTENT_JOB, intent); + } + + /** Arm `fill-window-expired` for an accepted intent (accept or a later amendment). */ + scheduleFillWindow(intent: Pick): void { + this.enqueue(FILL_WINDOW_EXPIRED_JOB, intent); + } + + private enqueue( + def: typeof EXPIRE_INTENT_JOB | typeof FILL_WINDOW_EXPIRED_JOB, + intent: Pick, + ): void { + if (!this.jobs) return; + try { + this.jobs.defineQueue(DEADLINE_QUEUE, { concurrency: 8 }); + void this.jobs + .enqueue( + def, + { intentId: intent.intentId, deadline: intent.deadline }, + { + delayMs: delayUntil(intent.deadline), + idempotencyKey: `${def.name}:${intent.intentId}:${intent.deadline}`, + }, + ) + .catch((err: unknown) => { + this.logger.error( + `[deadlines] failed to enqueue ${def.name} for ${intent.intentId}: ${ + err instanceof Error ? err.message : err + }`, + ); + }); + } catch (err) { + this.logger.error( + `[deadlines] failed to enqueue ${def.name} for ${intent.intentId}: ${ + err instanceof Error ? err.message : err + }`, + ); + } + } +} 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-sweeper.service.ts b/src/intents/intents-sweeper.service.ts index 98f355e..8b80968 100644 --- a/src/intents/intents-sweeper.service.ts +++ b/src/intents/intents-sweeper.service.ts @@ -1,4 +1,4 @@ -import { Injectable, Logger, OnModuleDestroy, OnModuleInit } from "@nestjs/common"; +import { Injectable, Logger, OnModuleDestroy, OnModuleInit, Optional } from "@nestjs/common"; import { IntentsService } from "./intents.service"; import { IntentsGateway } from "./intents.gateway"; import { SolversService } from "../solvers/solvers.service"; @@ -7,13 +7,18 @@ import { logger } from "../common/logger"; import { MetricsService } from "../metrics/metrics.service"; import { KillSwitchService } from "../killswitch/killswitch.service"; import { Intent } from "./intents.types"; +import { DeadlineJobData, DEADLINE_QUEUE, EXPIRE_INTENT_JOB, FILL_WINDOW_EXPIRED_JOB } from "./intents-deadline.jobs"; +import { JobsService } from "../jobs/jobs.service"; import { + AppConfig, CHAIN_FILL_WINDOW_DEFAULTS, DEFAULT_FILL_WINDOW_SECONDS, } from "../config/configuration"; import { LeaderElectionService, Singleton } from "../common/leader-election"; +import { ConfigService } from "@nestjs/config"; -const SWEEP_INTERVAL_MS = 30_000; +/** Low-frequency safety scan. Deadline jobs are the primary expiry path (issue #437). */ +const SAFETY_SWEEP_INTERVAL_MS = 300_000; /** Outcome of a single sweep cycle — returned so a manual trigger can log it. */ export interface SweepResult { @@ -38,9 +43,14 @@ export class IntentsSweeperService implements OnModuleInit, OnModuleDestroy { private readonly metricsService: MetricsService, private readonly killSwitch: KillSwitchService, private readonly leaderElection: LeaderElectionService, + @Optional() private readonly config?: ConfigService, + @Optional() private readonly jobs?: JobsService, ) {} onModuleInit() { + this.jobs?.defineQueue(DEADLINE_QUEUE, { concurrency: 8 }); + this.jobs?.process(EXPIRE_INTENT_JOB, (data) => this.handleExpireJob(data)); + this.jobs?.process(FILL_WINDOW_EXPIRED_JOB, (data) => this.handleFillWindowJob(data)); this.leaderElection.registerWorker("sweeper", (isLeader, _token) => { if (isLeader) { this.logger.log("[sweeper] became leader — starting interval"); @@ -59,10 +69,10 @@ export class IntentsSweeperService implements OnModuleInit, OnModuleDestroy { private startInterval(): void { if (this.interval) return; // already running this.interval = setInterval(() => { - this.sweep().catch((err) => { + this.sweep({ safety: true }).catch((err) => { logger.error(`[sweeper] sweep failed: ${err instanceof Error ? err.message : err}`); }); - }, SWEEP_INTERVAL_MS); + }, this.safetyIntervalMs()); } private stopInterval(): void { @@ -72,7 +82,7 @@ export class IntentsSweeperService implements OnModuleInit, OnModuleDestroy { } } - async sweep(): Promise { + async sweep(options?: { safety?: boolean }): Promise { const startMs = Date.now(); const now = Math.floor(startMs / 1000); let expiredCount = 0; @@ -80,22 +90,7 @@ export class IntentsSweeperService implements OnModuleInit, OnModuleDestroy { let extendedDeadlines = 0; for (const intent of await this.intentsService.getByState("open")) { - if (intent.deadline <= now) { - // Atomic guard: a concurrent user cancel() or solver accept() may have - // already transitioned this intent out of "open" — skip it if so. - const expired = await this.intentsService.expireIfOpen(intent.intentId); - if (!expired) continue; - // Audit trail (issue #62): system-driven expiration. - this.intentsService.appendAuditEntry( - intent.intentId, - "expired", - "system", - "deadline passed", - { deadline: intent.deadline, sweepedAt: now }, - ); - expiredCount++; - await this.intentsGateway.broadcast({ type: "intent_expired", intentId: intent.intentId }); - } + if (intent.deadline <= now && (await this.expireOpen(intent, now))) expiredCount++; } const durationMs = Date.now() - startMs; @@ -120,26 +115,13 @@ export class IntentsSweeperService implements OnModuleInit, OnModuleDestroy { // its window instead of slashing; the intent becomes fillable again on // resume. Evaluated per intent because a pause may be scoped to a single // chain or token. - const deadline = this.pausedFillDeadline(intent, now); - if (deadline !== null) { - const extended = await this.intentsService.extendDeadlineIfAccepted( - intent.intentId, - deadline, - ); - if (extended) { - extendedDeadlines++; - this.logger.warn( - `[sweeper] intent ${intent.intentId} fill is paused by a kill-switch — ` + - `slashing suppressed and deadline extended to ${deadline}`, - ); - } - continue; - } - - const slashed = await this.slashMissedFill(intent.intentId, intent.solver, now); - if (slashed) slashedCount++; + const outcome = await this.settleAccepted(intent, now); + if (outcome === "slashed") slashedCount++; + if (outcome === "extended") extendedDeadlines++; } + if (options?.safety) this.metricsService.recordSafetyCatch?.(expiredCount + slashedCount); + return { expiredCount, slashedCount, @@ -167,6 +149,65 @@ export class IntentsSweeperService implements OnModuleInit, OnModuleDestroy { return now + window; } + /** + * `expire-intent` handler. No-ops when the intent is no longer open or the + * stored deadline is not the one this job was armed for (amendment / stale job). + */ + async handleExpireJob(data: DeadlineJobData): Promise { + const intent = await this.intentsService.get(data.intentId); + if (!intent || intent.state !== "open" || intent.deadline !== data.deadline) return; + const now = Math.floor(Date.now() / 1000); + if (intent.deadline > now) return; + await this.expireOpen(intent, now); + } + + /** + * `fill-window-expired` handler. No-ops on a terminal state or a deadline + * that has since been moved. A kill-switch pause extends the window instead + * of slashing, and that extension arms a new job. + */ + async handleFillWindowJob(data: DeadlineJobData): Promise { + const intent = await this.intentsService.get(data.intentId); + if (!intent || intent.state !== "accepted" || intent.deadline !== data.deadline) return; + const now = Math.floor(Date.now() / 1000); + if (intent.deadline > now) return; + await this.settleAccepted(intent, now); + } + + private safetyIntervalMs(): number { + const configured = this.config?.get("safetySweepIntervalMs", { infer: true }); + return configured && configured > 0 ? configured : SAFETY_SWEEP_INTERVAL_MS; + } + + private async expireOpen(intent: Intent, now: number): Promise { + const expired = await this.intentsService.expireIfOpen(intent.intentId); + if (!expired) return false; + this.intentsService.appendAuditEntry( + intent.intentId, + "expired", + "system", + "deadline passed", + { deadline: intent.deadline, sweepedAt: now }, + ); + await this.intentsGateway.broadcast({ type: "intent_expired", intentId: intent.intentId }); + return true; + } + + private async settleAccepted(intent: Intent, now: number): Promise<"slashed" | "extended" | "skipped"> { + const deadline = this.pausedFillDeadline(intent, now); + if (deadline !== null) { + const extended = await this.intentsService.extendDeadlineIfAccepted(intent.intentId, deadline); + if (!extended) return "skipped"; + this.logger.warn( + `[sweeper] intent ${intent.intentId} fill is paused by a kill-switch — ` + + `slashing suppressed and deadline extended to ${deadline}`, + ); + return "extended"; + } + const slashed = await this.slashMissedFill(intent.intentId, intent.solver, now); + return slashed ? "slashed" : "skipped"; + } + /** * Issue #269 — safe, auditable manual sweep trigger (operator break-glass). * 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..4de9065 100644 --- a/src/intents/intents.module.ts +++ b/src/intents/intents.module.ts @@ -4,6 +4,7 @@ import { IntentsService } from "./intents.service"; import { IntentsController } from "./intents.controller"; import { IntentsGateway } from "./intents.gateway"; import { IntentsSweeperService } from "./intents-sweeper.service"; +import { IntentDeadlineScheduler } from "./intents-deadline.jobs"; import { IntentsMaintenanceJobs } from "./intents-maintenance.jobs"; import { INTENTS_REPOSITORY, InMemoryIntentsRepository } from "./intents.repository"; import { PrismaIntentsRepository } from "./prisma-intents.repository"; @@ -25,13 +26,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. @@ -54,6 +49,7 @@ import { GovernanceModule } from "../governance/governance.module"; IntentsGateway, backplaneHealthIndicator, IntentsSweeperService, + IntentDeadlineScheduler, IntentsMaintenanceJobs, // Note: EventIngestionService is provided by SorobanModule (imported above) // and exported from there — no re-declaration needed here. 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..327acd4 100644 --- a/src/intents/intents.service.ts +++ b/src/intents/intents.service.ts @@ -24,6 +24,7 @@ import { MetricsService } from "../metrics/metrics.service"; import { PrismaService } from "../prisma/prisma.service"; import { ProtocolParamsService } from "../governance/params.service"; import { FeatureFlagService } from "../flags/feature-flag.service"; +import { IntentDeadlineScheduler } from "./intents-deadline.jobs"; const TERMINAL_STATES: IntentState[] = ["filled", "cancelled", "expired", "slashed"]; @@ -109,6 +110,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,8 +130,8 @@ export class IntentsService { * this is always present. */ @Optional() private readonly metricsService?: MetricsService, - private readonly protocolParamsService: ProtocolParamsService, @Optional() private readonly flags?: FeatureFlagService, + @Optional() private readonly deadlines?: IntentDeadlineScheduler, ) {} /** @@ -267,6 +269,7 @@ export class IntentsService { } await this.repo.save(intent); + this.deadlines?.scheduleExpire(intent); // Creation is the entry edge of the funnel: the `vortex:intent:*` recording // rules count transitions *into* each state, so without this the intent // dashboard would start every conversion ratio from zero. `from_state` is @@ -513,7 +516,10 @@ export class IntentsService { const fillWindow = CHAIN_FILL_WINDOW_DEFAULTS[intent.srcChain] ?? DEFAULT_FILL_WINDOW_SECONDS; const updated = await this.repo.acceptIfOpen(id, solver, nowSec + fillWindow, nowSec); - if (updated !== null) this.countTransition("open", "accepted"); + if (updated !== null) { + this.countTransition("open", "accepted"); + this.deadlines?.scheduleFillWindow(updated); + } if (this.beginShadowObservation()) { this.observeAccept(updated ?? intent, solver, updated !== null); } @@ -533,9 +539,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 +657,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 +665,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" }), ]), ); } @@ -678,7 +682,9 @@ export class IntentsService { * or already has a later deadline. */ async extendDeadlineIfAccepted(id: string, newDeadline: number): Promise { - return this.repo.extendDeadlineIfAccepted(id, newDeadline); + const updated = await this.repo.extendDeadlineIfAccepted(id, newDeadline); + if (updated) this.deadlines?.scheduleFillWindow(updated); + return updated; } // --------------------------------------------------------------------------- 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/metrics/metrics.service.ts b/src/metrics/metrics.service.ts index 51d5047..1114bcb 100644 --- a/src/metrics/metrics.service.ts +++ b/src/metrics/metrics.service.ts @@ -76,6 +76,8 @@ export class MetricsService implements OnModuleInit { */ public readonly sweeperExpiredTotal: client.Counter; public readonly sweeperSweepDurationMs: client.Histogram; + /** Intents the low-frequency safety sweep expired or slashed. Steady state is ~0. */ + public readonly sweeperSafetyCaughtTotal: client.Counter; // ── SLO SLIs (issue #480) ───────────────────────────────────────────────── public readonly txConfirmationDuration: client.Histogram; @@ -155,6 +157,12 @@ export class MetricsService implements OnModuleInit { registers: [this.register], }); + this.sweeperSafetyCaughtTotal = new client.Counter({ + name: `${prefix}sweeper_safety_caught_total`, + help: "Intents expired or slashed by the low-frequency safety sweep (lost deadline jobs). Should stay near zero.", + registers: [this.register], + }); + // ── SLO SLIs (issue #480) ─────────────────────────────────────────────── this.intentCreateDuration = new client.Histogram({ name: `${prefix}intent_create_duration_seconds`, @@ -427,6 +435,11 @@ export class MetricsService implements OnModuleInit { this.sweeperSweepDurationMs.observe(durationMs); } + /** Items the safety sweep had to settle because a deadline job did not. */ + recordSafetyCatch(count: number): void { + if (count > 0) this.sweeperSafetyCaughtTotal.inc(count); + } + /** * Observe intent-create latency (SLO SLI, issue #480). * Call from the create path with handler duration in seconds. 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..a46d44a 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,25 @@ 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", + }, + safetySweepIntervalMs: 300_000, 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/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" } }