Skip to content
Merged
5 changes: 5 additions & 0 deletions CHANGELOG.md
Original file line number Diff line number Diff line change
Expand Up @@ -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).
Expand Down Expand Up @@ -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).
Expand Down
24 changes: 19 additions & 5 deletions README.md
Original file line number Diff line number Diff line change
Expand Up @@ -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.

---

Expand Down Expand Up @@ -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<string, string | null>, keyed by decimal stream-id string

// Protocol fee in basis points (e.g. 30 = 0.3%)
const feeBps = await client.factory.protocolFeeBps();
```
Expand All @@ -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
Expand Down
55 changes: 49 additions & 6 deletions docs/api.md
Original file line number Diff line number Diff line change
Expand Up @@ -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
Expand All @@ -36,7 +38,7 @@ new ConduitClient(config: ConduitConfig)

* `pauseStream(streamId: string) → Promise<string>` — equivalent to `client.streams.pause(streamId)`.
* `unpauseStream(streamId: string) → Promise<string>` — 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.

---

Expand Down Expand Up @@ -187,10 +189,14 @@ Resumes a paused stream. Paused duration is excluded from streaming time.

---

### `topUp(streamId, amount) → Promise<string>`
### `topUp(streamId, amount, signal?) → Promise<string>`

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`

Expand Down Expand Up @@ -329,32 +335,69 @@ endpoint; call `subscribe()` again to restart it.
## `client.factory`

### `streamCount() → Promise<bigint>`

### `streamAddress(id) → Promise<string | null>`

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<Map<string, string | null>>`

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<number>`

---

## `client.governor`

### `config() → Promise<GovernorConfig>`
### `getConfig(signal?) → Promise<GovernorConfig>`

```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`
Expand Down
14 changes: 11 additions & 3 deletions src/client.ts
Original file line number Diff line number Diff line change
Expand Up @@ -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()`.
Expand All @@ -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);
}

/**
Expand Down
86 changes: 86 additions & 0 deletions src/factory.ts
Original file line number Diff line number Diff line change
Expand Up @@ -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
Expand All @@ -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;
Expand Down Expand Up @@ -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<Map<string, string | null>> {
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<string>();
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<string, string | null>();
keys.forEach((key, i) => {
out.set(key, resolved[i] ?? null);
});
return out;
}

/**
* Whether `streamId` resolves to a deployed stream contract (#794).
*
Expand Down
Loading