diff --git a/CHANGELOG.md b/CHANGELOG.md index 2ce7fc1..cdd43db 100644 --- a/CHANGELOG.md +++ b/CHANGELOG.md @@ -7,6 +7,8 @@ All notable changes are documented here. Format based on [Keep a Changelog](http ### Added - `NETWORK_NAMES`, `EXPLORER_URLS`, and `NetworkType` provide shared human-readable Stellar network labels and Stellar Expert transaction, contract, and account URL bases (#832). - `TokenModule` exposes SEP-41 `allowance()` and `approve()` operations through `client.tokens`, and `@streamfi/react` now exports `useTokenAllowance()` for allowance verification and approval state (#851). +- `FactoryModule.streamAddresses(ids[], signal?, options?)` resolves a page of stream IDs to contract addresses in one call, returning a `Map` keyed by decimal stream-id string (`null` for ids the contract reports as not-found). Rendering a page of `streamsBySender()` results no longer costs one simulated RPC round trip per row on a cold cache: duplicate and string/bigint id forms are de-duplicated, only the cache-miss subset is fetched, and in-flight simulations are bounded by `options.maxConcurrency` (default 8) to match a default `stellar-rpc`'s 8 preflight workers. Every id is resolved through `streamAddress()`, so the address cache, its hit/miss counters and the negative-cache TTL are shared with single-id lookups; a resolution that throws propagates and is never cached as a not-found result. `StreamAddressesOptions` is exported from the package entry point (#783). +- `ConduitConfig.governorConfigCacheTtlMs` (default 30_000) bounds how long a `GovernorModule.getConfig()` read is reused. Protocol parameters only change when a governance proposal passes, so a dashboard polling `getConfig()` on an interval no longer pays a simulation per tick for data that is almost always unchanged. Concurrent misses share a single simulation, a failed simulation is never cached, each caller gets its own object so mutating a result cannot corrupt what the next caller sees, and `clearConfigCache()` forces the next call to re-simulate. Set the TTL to `0` to disable caching (#785). - `timeoutSignal(ms)` utility (exported from the package root and `/utils`) — a portable `AbortSignal` that aborts after `ms`, using the native `AbortSignal.timeout()` when available and falling back to `AbortController` + `setTimeout` (with `unref()` on Node) otherwise. Pass it as `signal` to any method that accepts one (#634). - `examples/quickstart.ts` — a runnable, end-to-end create -> accrue -> withdraw script on testnet, and the README Quickstart now mirrors it (#633). - `GraphQLIndexer.query()` now accepts optional `timeoutMs` (default 15s) and `signal` on `GraphQLQueryOptions` and wires a per-request `AbortController` into the underlying `fetch`, so a hung/slow indexer no longer leaves the caller's `await` pending forever. On timeout it rejects with a `IndexerTimeoutError` (endpoint + `timeoutMs` exposed); a caller-supplied `signal` surfaces the underlying `AbortError`. `IndexerTimeoutError` and `DEFAULT_INDEXER_TIMEOUT_MS` are exported from the package entry point (#569). @@ -34,6 +36,9 @@ All notable changes are documented here. Format based on [Keep a Changelog](http - `buildBatchTransactions()` (the RPC-prepared batch path) now simulates all operations in a batch concurrently instead of one at a time, cutting the wall-clock time of an N-operation batch from N sequential RPC round trips to one. ### Changed +- `ConduitClient.setWallet()` now propagates the new wallet to an already-constructed `FactoryModule`, whose read simulations are sourced from the active wallet. Previously the JSDoc documented `FactoryModule` as "not updated … does not hold a wallet reference and is unaffected by setWallet()", which contradicted the module's own `setWallet()`/`activeWallet` support and left a dApp that swapped wallets with a `client.factory` still simulating as the previous wallet. `FactoryModule.setWallet()` itself is unchanged and remains public API (#784). +- `StreamsModule.topUpStream()` is documented as the string-typed convenience wrapper over `topUp()` that it is, including which one new code should prefer. Both call the same contract method with the same validation; `topUpStream` additionally cannot forward an `AbortSignal` (#786). +- `StreamsModule.getStreamInfos({ maxConcurrency })` now treats a bound below 1 (or a non-finite one) as 1 instead of resolving to no results at all. The shared `mapWithConcurrency` helper and its default bound moved to `src/map-with-concurrency.ts` so the paged `list()` path and `FactoryModule.streamAddresses()` cannot drift apart (#783). - `u64ToScVal()` now rejects a non-integer `number` or a negative value with a clear `RangeError` naming the argument, and `estimateRequiredFee()` guards every `BigInt(...)` coercion (truncating a numeric input, catching an un-parseable one) and falls through to `fallbackStroops` — a non-conforming RPC response with a float `minResourceFee` no longer aborts a `create()` with a raw `RangeError` out of the fee-estimation path (#577). - `StreamsModule` now routes all signer-selection logic through its private `_signer()` helper instead of touching `config.signer` directly, removing dead code (#446). - CAIP-2→network mapping consolidated into a single exported `CAIP2_TO_NETWORK` constant shared by `ConduitClient`'s wallet network check and `WalletConnectAdapter`'s chain validation, so the two can never disagree (#445). diff --git a/README.md b/README.md index 2c86ebd..300d397 100644 --- a/README.md +++ b/README.md @@ -205,10 +205,11 @@ await walletAdapter.connect(); client.setWallet(walletAdapter); ``` -`setWallet()` only affects `client.streams` — `client.factory` and `client.governor` are -read-only and keep using `config.keypair` (if any) for simulation fee sourcing. It throws -`UnsupportedChainError` if the adapter's `chainId` doesn't match the network the client was -configured for. See [`setWallet`](docs/api.md#clientstreams) in the API reference for details. +`setWallet()` propagates to `client.streams` and to an already-constructed `client.factory` (whose +read simulations are sourced from the active wallet). `client.governor` is read-only and keeps +using `config.keypair` (if any) for simulation fee sourcing. It throws `UnsupportedChainError` if +the adapter's `chainId` doesn't match the network the client was configured for. See +[`setWallet`](docs/api.md#clientstreams) in the API reference for details. --- @@ -502,6 +503,13 @@ const count = await client.factory.streamCount(); // Stream address by ID const address = await client.factory.streamAddress(streamId); +// Addresses for a whole page, in one call — renders a stream list without +// one RPC round trip per row. Ids already resolved (or cached as +// not-found) cost nothing; only the cache-miss subset is fetched. +const { ids } = await client.factory.streamsBySender(sender, 0, 50); +const addresses = await client.factory.streamAddresses(ids); +// Map, keyed by decimal stream-id string + // Protocol fee in basis points (e.g. 30 = 0.3%) const feeBps = await client.factory.protocolFeeBps(); ``` @@ -513,17 +521,23 @@ const feeBps = await client.factory.protocolFeeBps(); Read protocol configuration: ```typescript -const config = await client.governor.config(); +const config = await client.governor.getConfig(); // Returns: // { // feeBps: number, // feeRecipient?: string, // minDurationSeconds: number, +// maxDurationSeconds: number, // maxRatePerSecond: bigint, // factoryAddress?: string, // } ``` +The result is reused for `governorConfigCacheTtlMs` (default 30s) — protocol parameters only +change when a governance proposal passes, so a polling dashboard no longer pays a simulation +per tick. Set it to `0` to re-simulate on every call, or call +`client.governor.clearConfigCache()` to force a refresh. + --- ## Error Handling diff --git a/docs/api.md b/docs/api.md index 3c391e7..7836896 100644 --- a/docs/api.md +++ b/docs/api.md @@ -22,6 +22,8 @@ new ConduitClient(config: ConduitConfig) | `wallet` | `WalletAdapter` | | — | | `fee` | `string` | | Explicit inclusion (bid) fee in stroops for submitted transactions. Takes precedence over `feeMultiplier`. Defaults to `BASE_FEE` (100 stroops) | | `feeMultiplier` | `number` | | Multiplier applied to `BASE_FEE` to compute the inclusion fee, e.g. `10` bids 10x the network minimum. Ignored when `fee` is set | +| `negativeCacheTtlMs` | `number` | | How long a not-found `streamAddress()` result is cached. Default 30_000 | +| `governorConfigCacheTtlMs` | `number` | | How long a `governor.getConfig()` read is reused. Default 30_000. `0` disables the cache | > **Inclusion fee.** Every submitted transaction previously used `BASE_FEE` > (the network minimum, 100 stroops) unconditionally, with no way to raise @@ -36,7 +38,7 @@ new ConduitClient(config: ConduitConfig) * `pauseStream(streamId: string) → Promise` — equivalent to `client.streams.pause(streamId)`. * `unpauseStream(streamId: string) → Promise` — equivalent to `client.streams.resume(streamId)`. -* `setWallet(wallet: WalletAdapter): void` — dynamically attach or change the active wallet adapter. Throws `UnsupportedChainError` if the wallet's `chainId` is on a different network than the client was configured for. See [Wallet Adapters](#wallet-adapters) below. Propagates to `client.streams` and `client.tokens`; `client.factory` and `client.governor` are read-only and unaffected. +* `setWallet(wallet: WalletAdapter): void` — dynamically attach or change the active wallet adapter. Throws `UnsupportedChainError` if the wallet's `chainId` is on a different network than the client was configured for. See [Wallet Adapters](#wallet-adapters) below. Propagates to `client.streams` and `client.tokens`, and to an already-constructed `client.factory` (whose read simulations are sourced from the active wallet). `client.governor` is read-only, uses `config.keypair` for simulation fee sourcing, and is unaffected. --- @@ -187,10 +189,14 @@ Resumes a paused stream. Paused duration is excluded from streaming time. --- -### `topUp(streamId, amount) → Promise` +### `topUp(streamId, amount, signal?) → Promise` Adds tokens to the stream balance. Extends effective stream duration. +`amount` is a `bigint` in stroops. For callers that already hold the amount as a string, +`topUpStream(streamId, amount)` is a thin string-typed wrapper that coerces it and delegates — +prefer `topUp` in new code, since it also accepts an `signal`. + **Requires:** `keypair` set (sender) **Throws:** `Error` (client-side) if `amount` is `<= 0n` — validated before any RPC round-trip; `ConduitError` with `contract: 'stream'` — `StreamErrorCode.StreamCancelled`, `.InvalidAmount` @@ -329,13 +335,42 @@ endpoint; call `subscribe()` again to restart it. ## `client.factory` ### `streamCount() → Promise` + ### `streamAddress(id) → Promise` Resolved (non-null) addresses are cached in-memory for the lifetime of the client, since a stream's contract address is fixed at creation and never changes. A `null` result (stream not -yet found) is not cached, so a later call for the same `id` will still hit the network. This -cache is what `StreamsModule` relies on to avoid re-resolving the same address on every -`get`/`withdraw`/`cancel`/`pause`/`resume`/`topUp`/`clawback` call and when paginating `list()`. +yet found) is cached for `ConduitConfig.negativeCacheTtlMs` (default 30s) so a polled list page +does not re-simulate every missing id on every refresh; call `clearAddressCache()` to drop it +earlier. This cache is what `StreamsModule` relies on to avoid re-resolving the same address on +every `get`/`withdraw`/`cancel`/`pause`/`resume`/`topUp`/`clawback` call and when paginating +`list()`. + +### `streamAddresses(ids, signal?, options?) → Promise>` + +Resolves a whole page of stream IDs in one call, instead of one simulated RPC round trip per id. + +```typescript +const { ids } = await client.factory.streamsBySender(sender, 0, 50); +const addresses = await client.factory.streamAddresses(ids); + +for (const [id, address] of addresses) { + // address is null when the id does not resolve to a deployed stream +} +``` + +| Param | Type | Notes | +|-------|------|-------| +| `ids` | `(bigint \| string)[]` | Duplicates — including the same id in both forms — are fetched once | +| `signal` | `AbortSignal?` | Rejects with an `AbortError` if aborted before or during resolution | +| `options.maxConcurrency` | `number?` | Max simulations in flight; default 8 (a default `stellar-rpc` serves `simulateTransaction` from 8 preflight workers). Values below 1 are treated as 1 | + +The returned `Map` is keyed by the decimal stream-id string and preserves first-seen input +order. Every id is resolved through `streamAddress()`, so the address cache, its hit/miss +counters and the negative-cache TTL above are shared with single-id lookups in both directions — +only the cache-miss subset costs a network call. A resolution that *fails* (as opposed to +resolving to "not found") rejects the whole call and is never cached as a not-found result; ids +that did resolve stay cached, so a retry re-fetches only the failures. ### `protocolFeeBps() → Promise` @@ -343,18 +378,26 @@ cache is what `StreamsModule` relies on to avoid re-resolving the same address o ## `client.governor` -### `config() → Promise` +### `getConfig(signal?) → Promise` ```typescript interface GovernorConfig { feeBps: number; feeRecipient?: string; minDurationSeconds: number; + maxDurationSeconds: number; maxRatePerSecond: bigint; factoryAddress?: string; } ``` +The result is reused for `ConduitConfig.governorConfigCacheTtlMs` (default 30s, ~6 ledgers at +the 5s close cadence), since protocol parameters only change when a governance proposal passes — +a dashboard polling `getConfig()` no longer pays a simulation per tick. Concurrent misses share +a single simulation, a failed simulation is never cached, and each caller receives its own object +so mutating a result cannot change what the next caller sees. Set `governorConfigCacheTtlMs: 0` +to re-simulate on every call, or call `clearConfigCache()` to force the next call to re-simulate. + --- ## `StreamInfo` diff --git a/src/client.ts b/src/client.ts index ada880c..83f9905 100644 --- a/src/client.ts +++ b/src/client.ts @@ -348,9 +348,11 @@ export class ConduitClient { * operations (create, withdraw, cancel, etc.) use the new wallet. * - {@link TokenModule}: Updated immediately — subsequent token approvals * use the new wallet. - * - {@link FactoryModule}: NOT updated — this module is read-only and - * uses `config.keypair` for simulation fee sourcing. It does not hold - * a wallet reference and is unaffected by `setWallet()`. + * - {@link FactoryModule}: Updated — its read simulations are sourced + * from the active wallet's public key, so a wallet swap re-resolves the + * simulation source instead of leaving it pinned to the previous + * wallet. A `client.factory` that has not been constructed yet picks + * the new wallet up from `config` when it is lazily built (#784). * - {@link GovernorModule}: NOT updated — this module is read-only and * uses `config.keypair` for simulation fee sourcing. It does not hold * a wallet reference and is unaffected by `setWallet()`. @@ -364,6 +366,12 @@ export class ConduitClient { this.config.wallet = wallet; this.streams.setWallet(wallet); this.tokens.setWallet(wallet); + // #784 — `FactoryModule` resolves its read-simulation source from a + // wallet adapter captured at construction, so it needs the swap too. + // Only when already constructed: a factory built later from the updated + // `config.wallet` above has nothing to catch up on. GovernorModule + // deliberately stays out — it holds no wallet at all. + this._factory?.setWallet(wallet); } /** diff --git a/src/factory.ts b/src/factory.ts index cebaf9c..4c35088 100644 --- a/src/factory.ts +++ b/src/factory.ts @@ -16,6 +16,7 @@ import { DEFAULT_RPC, } from './soroban.js'; import { SUPPORTED_NETWORKS, UnsupportedChainError } from './errors.js'; +import { mapWithConcurrency, DEFAULT_LIST_CONCURRENCY } from './map-with-concurrency.js'; /** * A `null` (not-found) `streamAddress` result is cached only briefly — a @@ -31,6 +32,16 @@ export interface FactoryStreamListResult { hasMore: boolean; } +/** Options for {@link FactoryModule.streamAddresses}. */ +export interface StreamAddressesOptions { + /** + * Maximum number of `stream_address` simulations in flight at once. + * Defaults to 8 — the number of preflight workers a default `stellar-rpc` + * serves `simulateTransaction` from. Values below 1 are treated as 1. + */ + maxConcurrency?: number; +} + export class FactoryModule { private readonly rpcUrl: string; private readonly passphrase: string; @@ -240,6 +251,81 @@ export class FactoryModule { this.negativeCacheExpiry.set(key, Date.now() + this._negativeCacheTtlMs); } + /** + * Resolve many stream IDs to their contract addresses in one call (#783). + * + * {@link streamAddress} costs one simulated RPC call per id, so rendering a + * page of `streamsBySender()` results costs one round trip per row on a + * cold cache. This resolves a whole page at once: the already-cached ids + * (positive hits, and negative hits still inside their + * `negativeCacheTtlMs` window) are served from memory and only the + * cache-miss subset is fetched, with at most `maxConcurrency` simulations + * in flight so the client does not outrun the RPC's preflight workers. + * + * Ids may be `bigint` or `string`, and duplicates (including the same id + * in both forms) are fetched once. The returned `Map` is keyed by the + * decimal id string and preserves first-seen input order; an id the + * contract reports as `None` maps to `null`, exactly as + * {@link streamAddress} returns. + * + * Because each id is resolved through {@link streamAddress}, the address + * cache, its hit/miss counters and the negative-cache TTL are shared with + * single-id lookups in both directions. + * + * @param ids - Stream IDs to resolve. + * @param signal - Abort signal; rejects with an `AbortError` if aborted + * before or during resolution. + * @param options - See {@link StreamAddressesOptions}. + * @returns A `Map` of decimal stream-id string to contract address, or + * `null` for ids that do not resolve to a deployed stream. + * @throws If any single resolution fails. Ids that did resolve stay + * cached, so a retry only re-fetches the ones that failed. A failed + * resolution is never recorded as a not-found result. + * + * @example + * ```typescript + * const { ids } = await client.factory.streamsBySender(sender, 0, 50); + * const addresses = await client.factory.streamAddresses(ids); + * for (const [id, address] of addresses) { + * console.log(id, address ?? 'not found'); + * } + * ``` + */ + async streamAddresses( + ids: (bigint | string)[], + signal?: AbortSignal, + options: StreamAddressesOptions = {}, + ): Promise> { + if (signal?.aborted) throw new DOMException("Aborted", "AbortError"); + + // De-duplicate first so a page that repeats an id (or mixes the string + // and bigint forms of one) never schedules two simulations for it. + const keys: string[] = []; + const seen = new Set(); + for (const id of ids) { + const key = BigInt(id).toString(); + if (!seen.has(key)) { + seen.add(key); + keys.push(key); + } + } + + // Cached ids resolve without a network call, so they simply occupy a + // slot in the pool for no RPC cost; keeping one code path means the + // batch method can never drift from the single-id one. + const resolved = await mapWithConcurrency( + keys, + options.maxConcurrency ?? DEFAULT_LIST_CONCURRENCY, + (key) => this.streamAddress(key, signal), + ); + + const out = new Map(); + keys.forEach((key, i) => { + out.set(key, resolved[i] ?? null); + }); + return out; + } + /** * Whether `streamId` resolves to a deployed stream contract (#794). * diff --git a/src/governor.ts b/src/governor.ts index 15365c6..f322fc0 100644 --- a/src/governor.ts +++ b/src/governor.ts @@ -5,6 +5,7 @@ import { Address, xdr } from '@stellar/stellar-sdk'; import type { ConduitConfig, GovernorConfig } from './types/index.js'; import { ZERO_ADDR } from './constants.js'; +import { coalesceAsync } from './coalesce-async.js'; import { buildContractCallTx, simulateReadOnly, @@ -16,6 +17,18 @@ import { } from './soroban.js'; import { SUPPORTED_NETWORKS, UnsupportedChainError } from './errors.js'; +/** + * How long a `config` read is reused before re-simulating (#785). + * + * Governance parameters change only when a proposal passes, so the window + * is a staleness budget rather than a correctness requirement. 30s is ~6 + * Stellar ledgers at the 5s close cadence, and matches the default + * `negativeCacheTtlMs` so the two caches in this SDK fail fresh at the same + * time. Override it with `ConduitConfig.governorConfigCacheTtlMs`, or set + * `0` to disable caching. + */ +const GOVERNOR_CONFIG_CACHE_TTL_MS = 30_000; + export class GovernorModule { private readonly rpcUrl: string; private readonly passphrase: string; @@ -23,6 +36,27 @@ export class GovernorModule { private readonly callerAddr: string; private readonly network: ConduitConfig['network']; + /** + * Cached `config` read and the timestamp it goes stale at (#785). A single + * entry is enough: a `GovernorModule` reads exactly one contract, so there + * is only ever one key. + */ + private configCache: { value: GovernorConfig; expiresAt: number } | undefined; + private readonly configCacheTtlMs: number; + /** + * In-flight `config` fetches, keyed by the constant `'config'`. Coalesces + * the burst of calls a TTL cache produces the moment its entry expires — + * without it, N concurrent callers all miss at once and N simulations go + * out. Reuses {@link coalesceAsync} (the same helper `soroban.ts` uses for + * token decimals), which evicts a rejection so the next caller retries. + * + * Entries are released as soon as the shared fetch settles, because + * `coalesceAsync` only evicts on rejection and a retained *fulfilled* + * promise would pin the first result past its TTL. An in-flight entry is + * therefore always a fetch that has not completed yet. + */ + private readonly inFlightConfig = new Map>(); + // Unlike FactoryModule (a hard prerequisite for virtually all StreamsModule // methods), GovernorModule is orthogonal to stream operations — a caller // using only client.streams shouldn't be forced to supply a @@ -40,9 +74,27 @@ export class GovernorModule { this.governorId = cfg.governorAddress; this.callerAddr = cfg.keypair?.publicKey() ?? ZERO_ADDR; this.network = cfg.network; + this.configCacheTtlMs = cfg.governorConfigCacheTtlMs ?? GOVERNOR_CONFIG_CACHE_TTL_MS; } - /** Fetch the current protocol config from the DripGovernor contract. */ + /** + * Fetch the current protocol config from the DripGovernor contract. + * + * The result is cached for `governorConfigCacheTtlMs` (default 30s) — + * governance parameters only change when a proposal passes, so a polling + * dashboard should not pay a simulation per tick for data that is almost + * always unchanged (#785). Concurrent misses share a single simulation, + * and a failed one is never cached. + * + * Each caller receives its own object, so mutating a result (e.g. holding + * a locally-adjusted copy) cannot change what the next caller sees. Call + * {@link clearConfigCache} to force the next call to re-simulate, or set + * `governorConfigCacheTtlMs: 0` to disable caching outright. + * + * @param signal - Abort signal; rejects with an `AbortError` if already + * aborted. The signal gates *this* call only — it never aborts an + * in-flight simulation that other concurrent callers are also awaiting. + */ async getConfig(signal?: AbortSignal): Promise { if (signal?.aborted) throw new DOMException("Aborted", "AbortError"); if (!this.governorId) { @@ -50,12 +102,39 @@ export class GovernorModule { `ConduitConfig.governorAddress is required (no default DripGovernor is known for network "${this.network}").`, ); } - const tx = await buildContractCallTx( - this.rpcUrl, this.passphrase, this.callerAddr, - this.governorId, 'config', [], - ); - const val = await simulateReadOnly(this.rpcUrl, this.passphrase, tx); - return parseGovernorConfig(val); + + const cached = this.configCache; + if (cached && Date.now() < cached.expiresAt) { + return { ...cached.value }; + } + + const value = await coalesceAsync(this.inFlightConfig, 'config', async () => { + const tx = await buildContractCallTx( + this.rpcUrl, this.passphrase, this.callerAddr, + this.governorId!, 'config', [], + ); + const val = await simulateReadOnly(this.rpcUrl, this.passphrase, tx); + const parsed = parseGovernorConfig(val); + // Stamped on completion, so the TTL bounds staleness from the moment + // the data was actually read rather than from when the call started. + // A TTL of 0 (or less) yields an entry that is never fresh again, so + // caching is effectively off without a separate code path. + this.configCache = { value: parsed, expiresAt: Date.now() + this.configCacheTtlMs }; + return parsed; + }); + + // `coalesceAsync` holds on to a *fulfilled* promise, which would pin the + // first result forever and defeat the TTL. Release the entry once the + // shared fetch has settled (callers already awaiting it keep their + // reference); a rejection was already evicted by `coalesceAsync` itself. + this.inFlightConfig.delete('config'); + + return { ...value }; + } + + /** Drop the cached `config` read so the next {@link getConfig} re-simulates. */ + clearConfigCache(): void { + this.configCache = undefined; } /** Fetch active governance proposals from the DripGovernor contract. */ diff --git a/src/index.ts b/src/index.ts index 699f0ad..e0150ec 100644 --- a/src/index.ts +++ b/src/index.ts @@ -168,6 +168,6 @@ export type { } from './module44.js'; export { FactoryModule } from './factory.js'; -export type { FactoryStreamListResult } from './factory.js'; +export type { FactoryStreamListResult, StreamAddressesOptions } from './factory.js'; export { GovernorModule } from './governor.js'; export type { GovernorProposal } from './governor.js'; diff --git a/src/map-with-concurrency.ts b/src/map-with-concurrency.ts new file mode 100644 index 0000000..3bbf1ee --- /dev/null +++ b/src/map-with-concurrency.ts @@ -0,0 +1,53 @@ +/** + * Bounded-concurrency `map` helper. + * + * Resolving N on-chain reads in a plain `Promise.all` fires N requests at + * the same instant; a Soroban RPC serves `simulateTransaction` from a fixed + * pool of preflight workers, so the extra fan-out turns into queue time and + * then rate-limit/timeouts. Running a small worker pool instead keeps the + * round trips overlapping without letting the client outrun the endpoint. + * + * Shared by `StreamsModule` (paged `list()` / `getStreamInfos()`) and + * `FactoryModule.streamAddresses()` (#783) so both use the same bound and + * the same result ordering. + */ + +/** + * Default ceiling on in-flight on-chain reads. + * + * Matches the 8 preflight workers a default `stellar-rpc` runs, so the + * default client concurrency is matched to the server's capacity rather + * than guessed at. + */ +export const DEFAULT_LIST_CONCURRENCY = 8; + +/** + * Runs `fn` over `items` with at most `concurrency` in-flight calls. + * Preserves result ordering to match a naive `Promise.all` fan-out, and + * rejects as soon as any call rejects (results already produced by other + * workers are still returned to their callers, they are just not collected). + * + * A `concurrency` below 1, or one that is not a finite number, is treated + * as 1 so a caller-supplied option can never silently resolve to no work + * at all. + */ +export async function mapWithConcurrency( + items: T[], + concurrency: number, + fn: (item: T) => Promise, +): Promise { + const results = new Array(items.length); + const limit = Math.max(1, Math.trunc(concurrency) || 1); + let index = 0; + + async function worker() { + while (index < items.length) { + const i = index++; + results[i] = await fn(items[i]!); + } + } + + const workers = Array.from({ length: Math.min(limit, items.length) }, () => worker()); + await Promise.all(workers); + return results; +} diff --git a/src/streams.ts b/src/streams.ts index 584857d..7bee9ba 100644 --- a/src/streams.ts +++ b/src/streams.ts @@ -51,6 +51,7 @@ import { import { buildBatchTransactions } from './batch-tx.js'; import type { BatchTransactionContext } from './batch-tx.js'; import { FactoryModule } from './factory.js'; +import { mapWithConcurrency, DEFAULT_LIST_CONCURRENCY } from './map-with-concurrency.js'; import { ConduitError, RateLimitError, @@ -67,32 +68,6 @@ import { * Tracks which v1-deprecated methods have already warned this session, so * repeated calls (e.g. in a hot loop) do not spam the console. */ -/** Default concurrency limit for bounded page-fetching (Issue #549). */ -const DEFAULT_LIST_CONCURRENCY = 8; - -/** - * Runs `fn` over `items` with at most `concurrency` in-flight calls. - * Preserves result ordering to match a naive `Promise.all` fan-out. - */ -async function mapWithConcurrency( - items: T[], - concurrency: number, - fn: (item: T) => Promise, -): Promise { - const results = new Array(items.length); - let index = 0; - - async function worker() { - while (index < items.length) { - const i = index++; - results[i] = await fn(items[i]!); - } - } - - const workers = Array.from({ length: Math.min(concurrency, items.length) }, () => worker()); - await Promise.all(workers); - return results; -} const _warnedDeprecations = new Set(); @@ -639,7 +614,12 @@ export class StreamsModule { return this._invoke(await this._resolveAddr(BigInt(streamId), signal), 'resume', [], signal); } - /** Deposit additional tokens into the stream (sender only). */ + /** + * Deposit additional tokens into the stream (sender only). + * + * The primary API: takes a `bigint` amount in stroops and an optional + * `signal`. See {@link topUpStream} for the string-typed wrapper. + */ async topUp(streamId: bigint | string, amount: bigint, signal?: AbortSignal): Promise { this._ensureCanMutate(); if (signal?.aborted) throw new DOMException('Aborted', 'AbortError'); @@ -670,7 +650,18 @@ export class StreamsModule { return this._invoke(await this._resolveAddr(BigInt(streamId)), 'force_cancel', []); } - /** Alias for topUp. */ + /** + * String-typed convenience wrapper over {@link topUp} — it coerces + * `amount` to a `bigint` and delegates, with no behaviour of its own. + * + * It exists for callers that already hold the amount as a string (form + * input, a `CreateStreamParams`-shaped value) and would otherwise have to + * convert before calling. **Prefer {@link topUp} in new code**: it is the + * primary method, takes the `bigint` amount the SDK uses for every other + * on-chain value, and accepts an `AbortSignal`, which this wrapper cannot + * forward. This is not a replacement for `topUp` and is not deprecated — + * both call the same contract method with the same validation. + */ async topUpStream(streamId: bigint | string, amount: bigint | string): Promise { return this.topUp(streamId, BigInt(amount)); } diff --git a/src/tests/client-set-wallet-propagation.test.ts b/src/tests/client-set-wallet-propagation.test.ts new file mode 100644 index 0000000..bd9aebf --- /dev/null +++ b/src/tests/client-set-wallet-propagation.test.ts @@ -0,0 +1,142 @@ +/** + * Tests for #784 — `ConduitClient.setWallet()` wallet propagation. + * + * `ConduitClient.setWallet()`'s "Wallet propagation contract" doc block + * claimed `FactoryModule` "does not hold a wallet reference and is + * unaffected by setWallet()". It does hold one, and does implement + * `setWallet()`, with async caller-address resolution. A dApp calling + * `client.setWallet(newWallet)` therefore left `client.factory`'s simulated + * caller address pinned to whatever it resolved before — the doc described + * an intentional design that the module's own code contradicted. + * + * These tests describe behaviour that does not exist on `main` — on main + * the factory's simulated source stays on the first wallet. + */ + +import { describe, it, expect, vi, beforeEach } from 'vitest'; +import type { ConduitConfig } from '../types/index.js'; +import type { WalletAdapter } from '../adapters/types.js'; + +// ── Hoisted mocks (same pattern as src/tests/bug-fixes-570-568.test.ts) ────── + +const { mockBuildTx, mockSimulate } = vi.hoisted(() => ({ + mockBuildTx: vi.fn(), + mockSimulate: vi.fn(), +})); + +vi.mock('../soroban.js', () => ({ + buildContractCallTx: mockBuildTx, + simulateReadOnly: mockSimulate, + scValToU64: (v: { u64: () => { toString: () => string } }) => BigInt(v.u64().toString()), + scValToU32: (v: { u32: () => number }) => v.u32(), + scValToI128: () => 0n, + resolveFee: () => '100', + NETWORK_PASSPHRASE: { + testnet: 'Test SDF Network ; September 2015', + mainnet: 'Public Global Stellar Network ; September 2015', + local: 'Standalone Network ; February 2017', + }, + DEFAULT_RPC: { + testnet: 'https://soroban-testnet.stellar.org', + mainnet: 'https://mainnet.sorobanrpc.com', + local: 'http://localhost:8000/soroban/rpc', + }, + catchNetworkError: (label: string, fn: () => unknown) => fn(), +})); + +vi.mock('@stellar/stellar-sdk', async () => { + const actual = await vi.importActual('@stellar/stellar-sdk'); + class MockAddress { + constructor(private readonly addr: string) {} + toScVal() { return actual.xdr.ScVal.scvVoid(); } + toString() { return this.addr; } + static fromScVal() { return new MockAddress('CADDRESS'); } + static fromString(s: string) { return new MockAddress(s); } + } + return { ...actual, Address: MockAddress }; +}); + +// ── Helpers ─────────────────────────────────────────────────────────────────── + +const FACTORY_ADDR = 'CCWAMYJME27OHTPKVSV252YRPXEO4BSKBHVLQ7ML3OWYNMB5RQEVHSM'; +const WALLET1_ADDR = 'GAAZI4TCR3TY5OJHCTJC2A4QSY6CJWJH5IAJTGKIN2ER7LBNVKOCCWN'; +const WALLET2_ADDR = 'GBBD47IF6LWK7P7MDEVSCWR7DPUWV3NY3DTQEVFL4NAT4AQH3ZLLFLA5'; + +function cfg(extra: Partial = {}): ConduitConfig { + return { network: 'testnet', factoryAddress: FACTORY_ADDR, rpcUrl: 'https://soroban-testnet.stellar.org', ...extra }; +} + +function stubWallet(pubkey: string): WalletAdapter { + return { + getPublicKey: vi.fn().mockResolvedValue(pubkey), + signTransaction: vi.fn((tx: unknown) => Promise.resolve(tx as never)), + }; +} + +function u64ScVal(n: bigint) { + return { switch: () => ({ name: 'scvU64' }), u64: () => ({ toString: () => n.toString() }) }; +} + +beforeEach(() => { + mockBuildTx.mockReset().mockResolvedValue({ _stub: 'tx' }); + mockSimulate.mockReset().mockResolvedValue(u64ScVal(0n)); +}); + +// ── Tests ───────────────────────────────────────────────────────────────────── + +describe('ConduitClient.setWallet() — factory propagation (#784)', () => { + it('re-resolves the factory read-simulation source on an already-constructed factory', async () => { + const { ConduitClient } = await import('../client.js'); + const client = new ConduitClient(cfg({ wallet: stubWallet(WALLET1_ADDR) })); + + // Force construction, then resolve once so the caller address is cached. + await client.factory.streamCount(); + client.setWallet(stubWallet(WALLET2_ADDR)); + await client.factory.streamCount(); + + expect(mockBuildTx.mock.calls[0]![2]).toBe(WALLET1_ADDR); + expect(mockBuildTx.mock.calls[1]![2]).toBe(WALLET2_ADDR); + }); + + it('applies a wallet set before the factory is ever constructed', async () => { + const { ConduitClient } = await import('../client.js'); + const client = new ConduitClient(cfg()); + + client.setWallet(stubWallet(WALLET2_ADDR)); + // The lazily constructed factory must pick the new wallet up from + // `config`, not from a stale construction-time snapshot. + await client.factory.streamCount(); + + expect(mockBuildTx.mock.calls[0]![2]).toBe(WALLET2_ADDR); + }); + + it('still propagates to StreamsModule', async () => { + const { ConduitClient } = await import('../client.js'); + const client = new ConduitClient(cfg({ wallet: stubWallet(WALLET1_ADDR) })); + const spy = vi.spyOn(client.streams, 'setWallet'); + + client.setWallet(stubWallet(WALLET2_ADDR)); + + expect(spy).toHaveBeenCalledTimes(1); + }); + + it('rejects a cross-network wallet before propagating to any module', async () => { + const { ConduitClient } = await import('../client.js'); + const client = new ConduitClient(cfg()); + const before = mockBuildTx.mock.calls.length; + const mainnetWallet: WalletAdapter = { + getPublicKey: vi.fn().mockResolvedValue(WALLET1_ADDR), + signTransaction: vi.fn((tx: unknown) => Promise.resolve(tx as never)), + chainId: 'stellar:mainnet', + }; + + expect(() => client.setWallet(mainnetWallet)).toThrow(/Unsupported network/); + await client.factory.streamCount(); + + // The rejected wallet never reached the factory: the source is still + // the ZERO_ADDR fallback, and the client config is unchanged. + const { ZERO_ADDR } = await import('../constants.js'); + expect(mockBuildTx.mock.calls.length).toBe(before + 1); + expect(mockBuildTx.mock.calls[before]![2]).toBe(ZERO_ADDR); + }); +}); diff --git a/src/tests/factory-stream-addresses.test.ts b/src/tests/factory-stream-addresses.test.ts new file mode 100644 index 0000000..e5eee30 --- /dev/null +++ b/src/tests/factory-stream-addresses.test.ts @@ -0,0 +1,373 @@ +/** + * Tests for #783 — `FactoryModule.streamAddresses(ids[])`. + * + * `streamAddress(id)` costs one simulated RPC call per id, so rendering a + * page of 50 streams from `streamsBySender()` costs 50 round trips on a cold + * cache. `streamAddresses()` resolves a whole page in one call, sharing + * `streamAddress()`'s cache, negative-cache TTL and abort semantics, and + * bounding how many simulations are in flight at once. + * + * These tests describe behaviour that does not exist on `main` — they fail + * there with "streamAddresses is not a function". + */ + +import { describe, it, expect, vi, beforeEach } from 'vitest'; +import { StrKey, xdr as _xdr } from '@stellar/stellar-sdk'; +import type { ConduitConfig } from '../types/index.js'; + +// ── Hoist mocks so they can be referenced inside vi.mock() factories ────────── + +const { mockBuildTx, mockSimulate } = vi.hoisted(() => ({ + mockBuildTx: vi.fn(), + mockSimulate: vi.fn(), +})); + +vi.mock('../soroban.js', () => ({ + buildContractCallTx: mockBuildTx, + simulateReadOnly: mockSimulate, + scValToU64: (v: { u64: () => { toString: () => string } }) => + BigInt(v.u64().toString()), + scValToI128: () => 0n, + scValToU32: (v: { u32: () => number }) => v.u32(), + NETWORK_PASSPHRASE: { + testnet: 'Test SDF Network ; September 2015', + mainnet: 'Public Global Stellar Network ; September 2015', + local: 'Standalone Network ; February 2017', + }, + DEFAULT_RPC: { + testnet: 'https://soroban-testnet.stellar.org', + mainnet: 'https://mainnet.sorobanrpc.com', + local: 'http://localhost:8000/soroban/rpc', + }, +})); + +// `Address.fromScVal` round-trips the real contract-id bytes so each id can +// resolve to its own distinct address (the upstream factory.test.ts mock +// returns a constant, which cannot distinguish one id's address from another's). +vi.mock('@stellar/stellar-sdk', async () => { + const actual = await vi.importActual('@stellar/stellar-sdk'); + + class MockAddress { + constructor(private readonly addr: string) {} + toScVal() { return actual.xdr.ScVal.scvVoid(); } + toString() { return this.addr; } + static fromScVal(v: { address: () => { contractId: () => Buffer } }) { + return new MockAddress(actual.StrKey.encodeContract(v.address().contractId())); + } + static fromString(s: string) { return new MockAddress(s); } + } + + return { ...actual, Address: MockAddress }; +}); + +// ── Helpers ─────────────────────────────────────────────────────────────────── + +const FACTORY_ADDR = 'CCWAMYJME27OHTPKVSV252YRPXEO4BSKBHVLQ7ML3OWYNMB5RQEVHSM'; + +function cfg(extra: Partial = {}): ConduitConfig { + return { + network: 'testnet', + factoryAddress: FACTORY_ADDR, + rpcUrl: 'https://soroban-testnet.stellar.org', + ...extra, + }; +} + +/** A distinct, valid contract address for each stream id used in the tests. */ +function addressFor(id: bigint): string { + // Zero-pad the decimal id into a 32-byte contract id, so every id used in + // these tests maps to its own valid C-address without any Number() cast. + return StrKey.encodeContract(Buffer.from(id.toString().padStart(64, '0'), 'hex')); +} + +/** `Option
` — a void ScVal is the contract's `None`. */ +function voidScVal(): _xdr.ScVal { + return _xdr.ScVal.scvVoid(); +} + +function addressScVal(id: bigint): _xdr.ScVal { + return _xdr.ScVal.scvAddress( + _xdr.ScAddress.scAddressTypeContract(StrKey.decodeContract(addressFor(id))), + ); +} + +/** + * Wire `buildContractCallTx` / `simulateReadOnly` so a `stream_address` + * simulation resolves per id according to `resolver`, and record every id the + * module actually sent to the network. + */ +function mockStreamAddress(resolver: (id: bigint) => _xdr.ScVal): { + simulated: bigint[]; +} { + const simulated: bigint[] = []; + mockBuildTx.mockImplementation( + async ( + _rpc: string, + _passphrase: string, + _source: string, + _contract: string, + _fn: string, + args?: _xdr.ScVal[], + ) => { + const id = BigInt(args![0]!.u64().toString()); + return { id }; + }, + ); + mockSimulate.mockImplementation( + async (_rpc: string, _passphrase: string, tx: { id: bigint }) => { + simulated.push(tx.id); + return resolver(tx.id); + }, + ); + return { simulated }; +} + +beforeEach(() => { + mockBuildTx.mockReset(); + mockSimulate.mockReset(); +}); + +// ── Tests ───────────────────────────────────────────────────────────────────── + +describe('FactoryModule.streamAddresses() (#783)', () => { + it('resolves every id in one call, keyed by decimal id string', async () => { + const { FactoryModule } = await import('../factory.js'); + const { simulated } = mockStreamAddress(id => addressScVal(id)); + + const map = await new FactoryModule(cfg()).streamAddresses([1n, 2n, 3n]); + + expect([...map.keys()]).toEqual(['1', '2', '3']); + expect(map.get('1')).toBe(addressFor(1n)); + expect(map.get('2')).toBe(addressFor(2n)); + expect(map.get('3')).toBe(addressFor(3n)); + expect(simulated).toHaveLength(3); + }); + + it('reports a not-found id as null rather than dropping the key', async () => { + const { FactoryModule } = await import('../factory.js'); + mockStreamAddress(id => (id === 2n ? voidScVal() : addressScVal(id))); + + const map = await new FactoryModule(cfg()).streamAddresses([1n, 2n, 3n]); + + expect(map.get('2')).toBeNull(); + expect(map.has('2')).toBe(true); + expect(map.size).toBe(3); + }); + + it('deduplicates repeated ids and collates string/bigint forms to one simulation', async () => { + const { FactoryModule } = await import('../factory.js'); + const { simulated } = mockStreamAddress(id => addressScVal(id)); + + const map = await new FactoryModule(cfg()).streamAddresses([7n, '7', 7n, 8n]); + + expect(map.size).toBe(2); + expect(map.get('7')).toBe(addressFor(7n)); + expect(simulated).toEqual([7n, 8n]); + }); + + it('preserves first-seen input order regardless of resolution order', async () => { + const { FactoryModule } = await import('../factory.js'); + // The first id resolves slowest, so a naive "push as answers arrive" + // implementation would emit a different order than the input. + mockBuildTx.mockImplementation( + async (_r: string, _p: string, _s: string, _c: string, _f: string, args?: _xdr.ScVal[]) => + ({ id: BigInt(args![0]!.u64().toString()) }), + ); + mockSimulate.mockImplementation(async (_r: string, _p: string, tx: { id: bigint }) => { + await new Promise(resolve => setTimeout(resolve, tx.id === 1n ? 30 : 0)); + return addressScVal(tx.id); + }); + + const map = await new FactoryModule(cfg()).streamAddresses([1n, 2n, 3n]); + + expect([...map.keys()]).toEqual(['1', '2', '3']); + }); + + it('accepts an empty id list without touching the network', async () => { + const { FactoryModule } = await import('../factory.js'); + mockStreamAddress(id => addressScVal(id)); + + const map = await new FactoryModule(cfg()).streamAddresses([]); + + expect(map.size).toBe(0); + expect(mockSimulate).not.toHaveBeenCalled(); + }); + + it('rejects with AbortError and skips the network when the signal is already aborted', async () => { + const { FactoryModule } = await import('../factory.js'); + mockStreamAddress(id => addressScVal(id)); + const controller = new AbortController(); + controller.abort(); + + await expect( + new FactoryModule(cfg()).streamAddresses([1n, 2n], controller.signal), + ).rejects.toMatchObject({ name: 'AbortError' }); + expect(mockSimulate).not.toHaveBeenCalled(); + }); + + it('rejects with AbortError when the signal aborts mid-flight', async () => { + const { FactoryModule } = await import('../factory.js'); + const controller = new AbortController(); + mockBuildTx.mockImplementation( + async (_r: string, _p: string, _s: string, _c: string, _f: string, args?: _xdr.ScVal[]) => + ({ id: BigInt(args![0]!.u64().toString()) }), + ); + mockSimulate.mockImplementation(async (_r: string, _p: string, tx: { id: bigint }) => { + if (tx.id === 1n) controller.abort(); + return addressScVal(tx.id); + }); + + // Serialised so the abort lands between two resolutions rather than + // before all of them were already scheduled. + await expect( + new FactoryModule(cfg()).streamAddresses([1n, 2n, 3n], controller.signal, { maxConcurrency: 1 }), + ).rejects.toMatchObject({ name: 'AbortError' }); + }); +}); + +describe('FactoryModule.streamAddresses() — cache sharing (#783)', () => { + it('only simulates the cache-miss subset when some ids are already resolved', async () => { + const { FactoryModule } = await import('../factory.js'); + const { simulated } = mockStreamAddress(id => addressScVal(id)); + const factory = new FactoryModule(cfg()); + + await factory.streamAddress(1n); + const map = await factory.streamAddresses([1n, 2n, 3n]); + + expect(map.get('1')).toBe(addressFor(1n)); + expect(simulated).toEqual([1n, 2n, 3n]); + }); + + it('reuses entries written by a previous streamAddresses() call', async () => { + const { FactoryModule } = await import('../factory.js'); + const { simulated } = mockStreamAddress(id => addressScVal(id)); + const factory = new FactoryModule(cfg()); + + await factory.streamAddresses([1n, 2n]); + const map = await factory.streamAddresses([1n, 2n]); + + expect(map.size).toBe(2); + expect(simulated).toHaveLength(2); + }); + + it('serves a negatively cached id from cache instead of re-simulating it', async () => { + const { FactoryModule } = await import('../factory.js'); + const { simulated } = mockStreamAddress(() => voidScVal()); + const factory = new FactoryModule(cfg({ negativeCacheTtlMs: 60_000 })); + + const first = await factory.streamAddresses([999n]); + const second = await factory.streamAddresses([999n, 998n]); + + expect(first.get('999')).toBeNull(); + expect(second.get('999')).toBeNull(); + // 999 came from the negative cache; only 998 hit the network. + expect(simulated).toEqual([999n, 998n]); + }); + + it('re-simulates a not-found id once its negative-cache TTL has expired', async () => { + const { FactoryModule } = await import('../factory.js'); + const { simulated } = mockStreamAddress(() => voidScVal()); + const factory = new FactoryModule(cfg({ negativeCacheTtlMs: 1 })); + vi.useFakeTimers(); + try { + vi.setSystemTime(0); + await factory.streamAddresses([999n]); + vi.setSystemTime(5_000); + await factory.streamAddresses([999n]); + } finally { + vi.useRealTimers(); + } + + expect(simulated).toEqual([999n, 999n]); + }); +}); + +describe('FactoryModule.streamAddresses() — bounded concurrency (#783)', () => { + /** Record the high-water mark of concurrently in-flight simulations. */ + function trackConcurrency(delayMs: number): { peak: () => number } { + let inFlight = 0; + let peak = 0; + mockBuildTx.mockImplementation( + async (_r: string, _p: string, _s: string, _c: string, _f: string, args?: _xdr.ScVal[]) => + ({ id: BigInt(args![0]!.u64().toString()) }), + ); + mockSimulate.mockImplementation(async (_r: string, _p: string, tx: { id: bigint }) => { + inFlight++; + peak = Math.max(peak, inFlight); + await new Promise(resolve => setTimeout(resolve, delayMs)); + inFlight--; + return addressScVal(tx.id); + }); + return { peak: () => peak }; + } + + it('never has more than maxConcurrency simulations in flight', async () => { + const { FactoryModule } = await import('../factory.js'); + const { peak } = trackConcurrency(5); + const ids = Array.from({ length: 24 }, (_, i) => BigInt(i + 1)); + + const map = await new FactoryModule(cfg()).streamAddresses(ids, undefined, { maxConcurrency: 3 }); + + expect(map.size).toBe(24); + expect(peak()).toBeLessThanOrEqual(3); + expect(peak()).toBeGreaterThan(1); + }); + + it('still resolves every id when maxConcurrency is below 1', async () => { + const { FactoryModule } = await import('../factory.js'); + const { peak } = trackConcurrency(1); + + const map = await new FactoryModule(cfg()).streamAddresses([1n, 2n, 3n], undefined, { maxConcurrency: 0 }); + + expect(map.size).toBe(3); + expect(peak()).toBe(1); + }); + + it('runs in parallel rather than sequentially at the default concurrency', async () => { + const { FactoryModule } = await import('../factory.js'); + const { peak } = trackConcurrency(20); + + const map = await new FactoryModule(cfg()).streamAddresses([1n, 2n, 3n, 4n]); + + expect(map.size).toBe(4); + expect(peak()).toBeGreaterThan(1); + }); +}); + +describe('FactoryModule.streamAddresses() — failure handling (#783)', () => { + it('rejects when an individual resolution fails, and does not cache it as not-found', async () => { + const { FactoryModule } = await import('../factory.js'); + const { simulated } = mockStreamAddress(id => { + if (id === 2n) throw new Error('rpc unavailable'); + return addressScVal(id); + }); + const factory = new FactoryModule(cfg()); + + await expect(factory.streamAddresses([1n, 2n, 3n])).rejects.toThrow('rpc unavailable'); + + // The failing id must not have poisoned the negative cache, so a retry + // re-simulates it instead of serving a bogus `null` for the next TTL. + mockStreamAddress(id => addressScVal(id)); + const retry = await factory.streamAddresses([2n]); + expect(retry.get('2')).toBe(addressFor(2n)); + expect(simulated).toContain(2n); + }); + + it('keeps the ids that did resolve cached, so a retry only re-fetches the failure', async () => { + const { FactoryModule } = await import('../factory.js'); + const failing = mockStreamAddress(id => { + if (id === 2n) throw new Error('rpc unavailable'); + return addressScVal(id); + }); + const factory = new FactoryModule(cfg()); + + await expect(factory.streamAddresses([1n, 2n])).rejects.toThrow('rpc unavailable'); + expect(failing.simulated).toContain(1n); + + const retry = mockStreamAddress(id => addressScVal(id)); + await factory.streamAddresses([1n, 2n]); + + // Id 1 was cached before the throw, so only id 2 is re-simulated. + expect(retry.simulated).toEqual([2n]); + }); +}); diff --git a/src/tests/factory.test.ts b/src/tests/factory.test.ts index ffbf6f3..7d86e50 100644 --- a/src/tests/factory.test.ts +++ b/src/tests/factory.test.ts @@ -302,22 +302,23 @@ describe('FactoryModule — hasStream() (#794)', () => { it('shares the streamAddress cache instead of re-hitting the network', async () => { const { FactoryModule } = await import('../factory.js'); - mockSimulate.mockResolvedValueOnce(makeU32ScVal(1)); + mockSimulate + .mockResolvedValueOnce(makeU32ScVal(1)) // id 5 resolves + .mockResolvedValueOnce(makeVoidScVal()); // id 999 is not found const factory = new FactoryModule(cfg()); const address = await factory.streamAddress(5n); const exists = await factory.hasStream('5'); const missing = await factory.hasStream(999n); - mockSimulate.mockResolvedValueOnce(makeVoidScVal()); - const missingResolved = await factory.hasStream(999n); + const missingAgain = await factory.hasStream(999n); expect(address).not.toBeNull(); expect(exists).toBe(true); expect(missing).toBe(false); - expect(missingResolved).toBe(false); - // Only three network hits total: id 5 resolved once, id 999 resolved - // twice (first miss cached negatively with a TTL, second miss explicit). - expect(mockSimulate).toHaveBeenCalledTimes(3); + expect(missingAgain).toBe(false); + // Two network hits total: id 5 resolved once, id 999 resolved once and + // was then served from the negative cache for the rest of its TTL. + expect(mockSimulate).toHaveBeenCalledTimes(2); }); it('propagates an already-aborted signal as AbortError without hitting the network', async () => { @@ -326,7 +327,9 @@ describe('FactoryModule — hasStream() (#794)', () => { const controller = new AbortController(); controller.abort(); - await expect(factory.hasStream(1n, controller.signal)).rejects.toThrow('AbortError'); + // 'AbortError' only appears in the exception's `name`; the message is + // 'Aborted', so matchObject on `name` rather than toThrow(string). + await expect(factory.hasStream(1n, controller.signal)).rejects.toMatchObject({ name: 'AbortError' }); expect(mockSimulate).not.toHaveBeenCalled(); }); }); diff --git a/src/tests/governor-config-cache.test.ts b/src/tests/governor-config-cache.test.ts new file mode 100644 index 0000000..cec4a3f --- /dev/null +++ b/src/tests/governor-config-cache.test.ts @@ -0,0 +1,280 @@ +/** + * Tests for #785 — `GovernorModule.getConfig()` TTL cache. + * + * `getConfig()` re-simulates the `config` contract call on every invocation + * with no caching, unlike `FactoryModule`'s address cache or `soroban.ts`'s + * token-decimals cache. Protocol parameters change rarely, so a dashboard + * polling `getConfig()` on an interval pays a full simulation round trip + * for data that is almost always unchanged. + * + * These tests describe behaviour that does not exist on `main` — they fail + * there because every call re-simulates. + */ + +import { describe, it, expect, vi, beforeEach } from 'vitest'; +import { Address, Keypair, StrKey, xdr as _xdr } from '@stellar/stellar-sdk'; +import type { ConduitConfig } from '../types/index.js'; + +// ── Hoist mocks so they can be referenced inside vi.mock() factories ────────── + +const { mockBuildTx, mockSimulate } = vi.hoisted(() => ({ + mockBuildTx: vi.fn(), + mockSimulate: vi.fn(), +})); + +vi.mock('../soroban.js', () => ({ + buildContractCallTx: mockBuildTx, + simulateReadOnly: mockSimulate, + scValToU64: (v: { u64: () => { toString: () => string } }) => BigInt(v.u64().toString()), + scValToU32: (v: { u32: () => number }) => v.u32(), + scValToI128: () => 0n, + NETWORK_PASSPHRASE: { + testnet: 'Test SDF Network ; September 2015', + mainnet: 'Public Global Stellar Network ; September 2015', + local: 'Standalone Network ; February 2017', + }, + DEFAULT_RPC: { + testnet: 'https://soroban-testnet.stellar.org', + mainnet: 'https://mainnet.sorobanrpc.com', + local: 'http://localhost:8000/soroban/rpc', + }, +})); + +// ── Helpers ─────────────────────────────────────────────────────────────────── + +const GOVERNOR_ADDR = 'CBQHNAXSI55GX2GN6D67GK7BHVPSLJUGZQEU7WJ5LKR5PNUCGLIMAO4K'; +const FEE_RECIPIENT = Keypair.random().publicKey(); +const FACTORY_ADDR = StrKey.encodeContract(Buffer.alloc(32, 9)); + +function cfg(extra: Partial = {}): ConduitConfig { + return { + network: 'testnet', + governorAddress: GOVERNOR_ADDR, + rpcUrl: 'https://soroban-testnet.stellar.org', + ...extra, + }; +} + +function scvMap(entries: Record): _xdr.ScVal { + return _xdr.ScVal.scvMap( + Object.entries(entries).map(([k, v]) => + new _xdr.ScMapEntry({ key: _xdr.ScVal.scvSymbol(k), val: v }), + ), + ); +} + +function u64(n: bigint) { return _xdr.ScVal.scvU64(_xdr.Uint64.fromString(n.toString())); } +function u32(n: number) { return _xdr.ScVal.scvU32(n); } + +function i128(n: bigint) { + const lo = n & 0xffffffffffffffffn; + const hi = n >> 64n; + return _xdr.ScVal.scvI128( + new _xdr.Int128Parts({ + hi: _xdr.Int64.fromString(hi.toString()), + lo: _xdr.Uint64.fromString(lo.toString()), + }), + ); +} + +/** A full, representative protocol config response. */ +function configScVal(feeBps = 30) { + return scvMap({ + fee_bps: u32(feeBps), + fee_recipient: new Address(FEE_RECIPIENT).toScVal(), + min_duration_seconds: u64(3_600n), + max_duration_seconds: u64(31_536_000n), + max_rate_per_second: i128(1_000_000_000_000_000n), + factory_address: new Address(FACTORY_ADDR).toScVal(), + }); +} + +beforeEach(() => { + mockBuildTx.mockReset().mockResolvedValue({ _stub: 'tx' }); + mockSimulate.mockReset(); +}); + +// ── Tests ───────────────────────────────────────────────────────────────────── + +describe('GovernorModule.getConfig() TTL cache (#785)', () => { + it('serves repeated calls within the TTL from cache (one simulation total)', async () => { + const { GovernorModule } = await import('../governor.js'); + mockSimulate.mockResolvedValue(configScVal()); + + const mod = new GovernorModule(cfg()); + const first = await mod.getConfig(); + const second = await mod.getConfig(); + const third = await mod.getConfig(); + + expect(first.feeBps).toBe(30); + expect(second.feeBps).toBe(30); + expect(third.feeBps).toBe(30); + expect(mockSimulate).toHaveBeenCalledTimes(1); + }); + + it('re-simulates once the TTL has expired and returns the new value', async () => { + const { GovernorModule } = await import('../governor.js'); + mockSimulate.mockResolvedValueOnce(configScVal(30)).mockResolvedValueOnce(configScVal(50)); + const mod = new GovernorModule(cfg({ governorConfigCacheTtlMs: 1_000 })); + + vi.useFakeTimers(); + try { + vi.setSystemTime(0); + const before = await mod.getConfig(); + expect(before.feeBps).toBe(30); + + // Still inside the TTL: the cached value wins even though the chain has + // already moved on. + vi.setSystemTime(500); + expect((await mod.getConfig()).feeBps).toBe(30); + expect(mockSimulate).toHaveBeenCalledTimes(1); + + // Past the TTL: a fresh simulation picks up the new parameter. + vi.setSystemTime(1_500); + expect((await mod.getConfig()).feeBps).toBe(50); + expect(mockSimulate).toHaveBeenCalledTimes(2); + } finally { + vi.useRealTimers(); + } + }); + + it('honours a custom governorConfigCacheTtlMs', async () => { + const { GovernorModule } = await import('../governor.js'); + mockSimulate.mockResolvedValue(configScVal()); + + const mod = new GovernorModule(cfg({ governorConfigCacheTtlMs: 60_000 })); + vi.useFakeTimers(); + try { + vi.setSystemTime(0); + await mod.getConfig(); + vi.setSystemTime(59_000); + await mod.getConfig(); + expect(mockSimulate).toHaveBeenCalledTimes(1); + vi.setSystemTime(61_000); + await mod.getConfig(); + expect(mockSimulate).toHaveBeenCalledTimes(2); + } finally { + vi.useRealTimers(); + } + }); + + it('re-simulates on every call when the TTL is 0 (caching disabled)', async () => { + const { GovernorModule } = await import('../governor.js'); + mockSimulate.mockResolvedValue(configScVal()); + + const mod = new GovernorModule(cfg({ governorConfigCacheTtlMs: 0 })); + await mod.getConfig(); + await mod.getConfig(); + await mod.getConfig(); + + expect(mockSimulate).toHaveBeenCalledTimes(3); + }); + + it('coalesces concurrent misses into a single simulation', async () => { + const { GovernorModule } = await import('../governor.js'); + mockSimulate.mockImplementation( + () => new Promise(resolve => setTimeout(() => resolve(configScVal()), 10)), + ); + + const mod = new GovernorModule(cfg()); + const results = await Promise.all([mod.getConfig(), mod.getConfig(), mod.getConfig()]); + + for (const r of results) expect(r.feeBps).toBe(30); + expect(mockSimulate).toHaveBeenCalledTimes(1); + }); + + it('does not cache a failed simulation, so the next call retries', async () => { + const { GovernorModule } = await import('../governor.js'); + mockSimulate + .mockRejectedValueOnce(new Error('rpc unavailable')) + .mockResolvedValue(configScVal()); + + const mod = new GovernorModule(cfg()); + await expect(mod.getConfig()).rejects.toThrow('rpc unavailable'); + const retried = await mod.getConfig(); + + expect(retried.feeBps).toBe(30); + expect(mockSimulate).toHaveBeenCalledTimes(2); + }); + + it('rejects every concurrent caller when the shared simulation fails', async () => { + const { GovernorModule } = await import('../governor.js'); + mockSimulate.mockImplementation( + () => new Promise((_resolve, reject) => setTimeout(() => reject(new Error('rpc unavailable')), 5)), + ); + + const mod = new GovernorModule(cfg()); + const results = await Promise.allSettled([mod.getConfig(), mod.getConfig()]); + + expect(results.every(r => r.status === 'rejected')).toBe(true); + expect(mockSimulate).toHaveBeenCalledTimes(1); + }); + + it('clearConfigCache() forces the next call to re-simulate', async () => { + const { GovernorModule } = await import('../governor.js'); + mockSimulate.mockResolvedValue(configScVal()); + + const mod = new GovernorModule(cfg()); + await mod.getConfig(); + mod.clearConfigCache(); + await mod.getConfig(); + + expect(mockSimulate).toHaveBeenCalledTimes(2); + }); + + it('hands each caller its own object, so mutating a result cannot corrupt the cache', async () => { + const { GovernorModule } = await import('../governor.js'); + mockSimulate.mockResolvedValue(configScVal()); + + const mod = new GovernorModule(cfg()); + const first = await mod.getConfig(); + const second = await mod.getConfig(); + + expect(first).not.toBe(second); + expect(first).toEqual(second); + + // A caller scribbling on a result must not change what the next caller + // sees — otherwise one component's local edit silently rewrites the + // protocol config for every other consumer in the process. + first.feeBps = 999; + first.minDurationSeconds = 1; + + const third = await mod.getConfig(); + expect(third.feeBps).toBe(30); + expect(third.minDurationSeconds).toBe(3_600); + expect(mockSimulate).toHaveBeenCalledTimes(1); + }); + + it('still throws the missing-governorAddress error before touching the network', async () => { + const { GovernorModule } = await import('../governor.js'); + const mod = new GovernorModule({ network: 'testnet', rpcUrl: 'https://soroban-testnet.stellar.org' }); + + await expect(mod.getConfig()).rejects.toThrow(/governorAddress is required/); + expect(mockSimulate).not.toHaveBeenCalled(); + }); + + it('rejects an already-aborted signal before consulting or populating the cache', async () => { + const { GovernorModule } = await import('../governor.js'); + mockSimulate.mockResolvedValue(configScVal()); + const mod = new GovernorModule(cfg()); + await mod.getConfig(); + + const controller = new AbortController(); + controller.abort(); + await expect(mod.getConfig(controller.signal)).rejects.toMatchObject({ name: 'AbortError' }); + + // Still one simulation: the abort neither served a value nor poisoned + // the cache with a rejected one. + expect(mockSimulate).toHaveBeenCalledTimes(1); + }); + + it('keeps caches separate per module instance', async () => { + const { GovernorModule } = await import('../governor.js'); + mockSimulate.mockResolvedValue(configScVal()); + + await new GovernorModule(cfg()).getConfig(); + await new GovernorModule(cfg()).getConfig(); + + expect(mockSimulate).toHaveBeenCalledTimes(2); + }); +}); diff --git a/src/tests/governor.test.ts b/src/tests/governor.test.ts index de2aad5..630967f 100644 --- a/src/tests/governor.test.ts +++ b/src/tests/governor.test.ts @@ -83,7 +83,9 @@ function i128(n: bigint) { // ── Tests ───────────────────────────────────────────────────────────────────── beforeEach(() => { - mockBuildTx.mockResolvedValue({ _stub: 'tx' }); + // mockBuildTx needs a full reset, not just a re-resolve: an earlier test's + // calls would otherwise leak into a later `not.toHaveBeenCalled()`. + mockBuildTx.mockReset().mockResolvedValue({ _stub: 'tx' }); mockSimulate.mockReset(); }); @@ -147,7 +149,9 @@ describe('GovernorModule — getConfig()', () => { const controller = new AbortController(); controller.abort(); - await expect(new GovernorModule(cfg()).getConfig(controller.signal)).rejects.toThrow('AbortError'); + // 'AbortError' only appears in the exception's `name`; the message is + // 'Aborted', so matchObject on `name` rather than toThrow(string). + await expect(new GovernorModule(cfg()).getConfig(controller.signal)).rejects.toMatchObject({ name: 'AbortError' }); expect(mockBuildTx).not.toHaveBeenCalled(); expect(mockSimulate).not.toHaveBeenCalled(); }); diff --git a/src/types/index.ts b/src/types/index.ts index fb4499c..62bd08a 100644 --- a/src/types/index.ts +++ b/src/types/index.ts @@ -56,6 +56,14 @@ export interface ConduitConfig { * retrying. Default 30_000 (30 seconds). */ negativeCacheTtlMs?: number; + /** + * How long a `GovernorModule.getConfig()` read is reused, in milliseconds. + * Protocol parameters only change when a governance proposal passes, so a + * short window removes the need to re-simulate on every poll. Default + * 30_000 (30 seconds, ~6 ledgers). Set `0` to disable the cache and + * re-simulate on every call. + */ + governorConfigCacheTtlMs?: number; } export interface StreamInfo {