Compare commits
3 commits
8b7098a4b0
...
762b0def5c
| Author | SHA1 | Date | |
|---|---|---|---|
| 762b0def5c | |||
| 3064acf217 | |||
| 2a64b42cde |
6 changed files with 301 additions and 42 deletions
|
|
@ -56,6 +56,17 @@ export interface OperatorConfigServiceConfig {
|
||||||
export interface OperatorConfigService {
|
export interface OperatorConfigService {
|
||||||
/** Unsubscribe from operator events and free resources. */
|
/** Unsubscribe from operator events and free resources. */
|
||||||
stop(): void
|
stop(): void
|
||||||
|
/**
|
||||||
|
* Republish the current cassette state (kind-30078, replaceable). Call after
|
||||||
|
* a dispense and on a cassette reload so the operator's view tracks reality.
|
||||||
|
* Best-effort — logs and swallows errors.
|
||||||
|
*/
|
||||||
|
publishCassettesState(): Promise<void>
|
||||||
|
}
|
||||||
|
|
||||||
|
const NOOP_SERVICE: OperatorConfigService = {
|
||||||
|
stop: () => {},
|
||||||
|
publishCassettesState: async () => {},
|
||||||
}
|
}
|
||||||
|
|
||||||
export async function startOperatorConfigService(
|
export async function startOperatorConfigService(
|
||||||
|
|
@ -63,11 +74,11 @@ export async function startOperatorConfigService(
|
||||||
): Promise<OperatorConfigService> {
|
): Promise<OperatorConfigService> {
|
||||||
if (cfg.operatorPubkeys.length === 0) {
|
if (cfg.operatorPubkeys.length === 0) {
|
||||||
console.log('[OperatorConfig] No operator pubkeys configured — service disabled')
|
console.log('[OperatorConfig] No operator pubkeys configured — service disabled')
|
||||||
return { stop: () => {} }
|
return NOOP_SERVICE
|
||||||
}
|
}
|
||||||
if (!isElectron || !window.electronAPI) {
|
if (!isElectron || !window.electronAPI) {
|
||||||
console.log('[OperatorConfig] Not in Electron — service disabled (browser dev mode)')
|
console.log('[OperatorConfig] Not in Electron — service disabled (browser dev mode)')
|
||||||
return { stop: () => {} }
|
return NOOP_SERVICE
|
||||||
}
|
}
|
||||||
const api = window.electronAPI
|
const api = window.electronAPI
|
||||||
const machineId = cfg.machineId ?? cfg.signer.pubkey
|
const machineId = cfg.machineId ?? cfg.signer.pubkey
|
||||||
|
|
@ -103,6 +114,12 @@ export async function startOperatorConfigService(
|
||||||
|
|
||||||
return {
|
return {
|
||||||
stop: () => cfg.nostrClient.unsubscribe(subscriptionId),
|
stop: () => cfg.nostrClient.unsubscribe(subscriptionId),
|
||||||
|
publishCassettesState: () =>
|
||||||
|
publishCassettesState(cfg, api, machineId)
|
||||||
|
.then(() => {})
|
||||||
|
.catch((err) => {
|
||||||
|
console.warn('[OperatorConfig] cassettes-state republish failed:', err)
|
||||||
|
}),
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
|
|
@ -193,29 +210,35 @@ async function handleOperatorConfigEvent(
|
||||||
console.log(
|
console.log(
|
||||||
`[OperatorConfig] Applied — created_at=${event.created_at}, positions=${Object.keys(parsed.positions).join(',')}`
|
`[OperatorConfig] Applied — created_at=${event.created_at}, positions=${Object.keys(parsed.positions).join(',')}`
|
||||||
)
|
)
|
||||||
|
|
||||||
|
// Republish our resulting cassette state so the operator's view reflects the
|
||||||
|
// applied config (the "on cassette reload" case). Different d-tag from the
|
||||||
|
// operator's config event, so no echo loop. Best-effort.
|
||||||
|
const machineId = cfg.machineId ?? cfg.signer.pubkey
|
||||||
|
await publishCassettesState(cfg, api, machineId).catch((err) =>
|
||||||
|
console.warn('[OperatorConfig] post-apply cassettes-state republish failed:', err)
|
||||||
|
)
|
||||||
}
|
}
|
||||||
|
|
||||||
async function maybePublishBootstrap(
|
/**
|
||||||
|
* Publish the ATM's current cassette state as a replaceable kind-30078 event
|
||||||
|
* (`bitspire-cassettes-state:<machineId>`), NIP-44-encrypted to the operator.
|
||||||
|
* Replaceable → latest wins; the operator consumes every update. Call after a
|
||||||
|
* dispense and on a cassette reload so the operator view tracks reality, not
|
||||||
|
* the frozen bootstrap snapshot (coord 2026-06-21 / lamassu-next#56).
|
||||||
|
*
|
||||||
|
* NOT gated on the bootstrap flag — this is the live update. Returns whether an
|
||||||
|
* event was published (false when there are no cassettes / no operator).
|
||||||
|
*/
|
||||||
|
async function publishCassettesState(
|
||||||
cfg: OperatorConfigServiceConfig,
|
cfg: OperatorConfigServiceConfig,
|
||||||
api: NonNullable<typeof window.electronAPI>,
|
api: NonNullable<typeof window.electronAPI>,
|
||||||
machineId: string
|
machineId: string
|
||||||
): Promise<void> {
|
): Promise<boolean> {
|
||||||
const already = await api.getBootstrapPublishedAt()
|
|
||||||
if (already !== null) {
|
|
||||||
console.log('[OperatorConfig] Bootstrap already published at unix', already)
|
|
||||||
return
|
|
||||||
}
|
|
||||||
const cassettes = await api.loadCassettes()
|
const cassettes = await api.loadCassettes()
|
||||||
if (cassettes.length === 0) {
|
if (cassettes.length === 0) return false
|
||||||
console.log('[OperatorConfig] state.db.cassettes empty — skipping bootstrap')
|
|
||||||
return
|
|
||||||
}
|
|
||||||
|
|
||||||
const operatorPubkey = cfg.operatorPubkeys[0]
|
const operatorPubkey = cfg.operatorPubkeys[0]
|
||||||
if (!operatorPubkey) {
|
if (!operatorPubkey) return false
|
||||||
console.log('[OperatorConfig] No operator pubkey — skipping bootstrap')
|
|
||||||
return
|
|
||||||
}
|
|
||||||
|
|
||||||
const positions: Record<string, { denomination: number; count: number }> = {}
|
const positions: Record<string, { denomination: number; count: number }> = {}
|
||||||
for (const c of cassettes) {
|
for (const c of cassettes) {
|
||||||
|
|
@ -235,6 +258,30 @@ async function maybePublishBootstrap(
|
||||||
})
|
})
|
||||||
|
|
||||||
await cfg.nostrClient.publish(event)
|
await cfg.nostrClient.publish(event)
|
||||||
await api.markBootstrapPublished(Math.floor(Date.now() / 1000))
|
console.log('[OperatorConfig] cassettes-state published:', { dTag, eventId: event.id })
|
||||||
console.log('[OperatorConfig] Bootstrap hello-event published:', { dTag, eventId: event.id })
|
return true
|
||||||
|
}
|
||||||
|
|
||||||
|
/**
|
||||||
|
* First-boot hello: publish the cassette state once and mark the gate. The
|
||||||
|
* gate (lamassu-next#56) prevents re-emitting the *bootstrap* on every boot;
|
||||||
|
* live updates after dispenses go through `publishCassettesState` directly.
|
||||||
|
*/
|
||||||
|
async function maybePublishBootstrap(
|
||||||
|
cfg: OperatorConfigServiceConfig,
|
||||||
|
api: NonNullable<typeof window.electronAPI>,
|
||||||
|
machineId: string
|
||||||
|
): Promise<void> {
|
||||||
|
const already = await api.getBootstrapPublishedAt()
|
||||||
|
if (already !== null) {
|
||||||
|
console.log('[OperatorConfig] Bootstrap already published at unix', already)
|
||||||
|
return
|
||||||
|
}
|
||||||
|
const published = await publishCassettesState(cfg, api, machineId)
|
||||||
|
if (published) {
|
||||||
|
await api.markBootstrapPublished(Math.floor(Date.now() / 1000))
|
||||||
|
console.log('[OperatorConfig] Bootstrap hello-event published')
|
||||||
|
} else {
|
||||||
|
console.log('[OperatorConfig] No cassettes/operator — skipping bootstrap')
|
||||||
|
}
|
||||||
}
|
}
|
||||||
|
|
|
||||||
|
|
@ -521,7 +521,10 @@ export const useAtmStore = defineStore('atm', () => {
|
||||||
bills,
|
bills,
|
||||||
cassettes: dr?.cassettes,
|
cassettes: dr?.cassettes,
|
||||||
error: dr?.error ?? ctx.error,
|
error: dr?.error ?? ctx.error,
|
||||||
}).then(() => reloadPersistedInventory())
|
})
|
||||||
|
.then(() => reloadPersistedInventory())
|
||||||
|
// Republish cassette state — a partial dispense changed counts.
|
||||||
|
.then(() => operatorConfigSvc?.publishCassettesState())
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
|
|
@ -551,7 +554,11 @@ export const useAtmStore = defineStore('atm', () => {
|
||||||
bills,
|
bills,
|
||||||
cassettes: dr?.cassettes,
|
cassettes: dr?.cassettes,
|
||||||
error: dr?.error,
|
error: dr?.error,
|
||||||
}).then(() => reloadPersistedInventory())
|
})
|
||||||
|
.then(() => reloadPersistedInventory())
|
||||||
|
// Republish cassette state after a cash-out dispense (counts
|
||||||
|
// decremented); harmless no-op echo for a cash-in complete.
|
||||||
|
.then(() => (isCashInTx ? undefined : operatorConfigSvc?.publishCassettesState()))
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
|
|
|
||||||
96
packages/lnbits/src/__tests__/retry.test.ts
Normal file
96
packages/lnbits/src/__tests__/retry.test.ts
Normal file
|
|
@ -0,0 +1,96 @@
|
||||||
|
import { describe, it, expect, vi } from 'vitest'
|
||||||
|
import { withRetry } from '../retry.js'
|
||||||
|
import { LnbitsErrorCode, LnbitsRpcError } from '../error-codes.js'
|
||||||
|
|
||||||
|
const noSleep = () => Promise.resolve()
|
||||||
|
|
||||||
|
function rpcErr(code: LnbitsErrorCode): LnbitsRpcError {
|
||||||
|
return new LnbitsRpcError({ code, rpcName: 'get_wallet', requestId: 'r1' })
|
||||||
|
}
|
||||||
|
|
||||||
|
/** A fn that throws `err` the first `failTimes` calls, then returns `value`. */
|
||||||
|
function failingFn<T>(failTimes: number, err: unknown, value: T): { fn: () => Promise<T>; calls: () => number } {
|
||||||
|
let calls = 0
|
||||||
|
return {
|
||||||
|
fn: async () => {
|
||||||
|
calls++
|
||||||
|
if (calls <= failTimes) throw err
|
||||||
|
return value
|
||||||
|
},
|
||||||
|
calls: () => calls,
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
describe('withRetry', () => {
|
||||||
|
it('returns immediately on success (one call)', async () => {
|
||||||
|
const { fn, calls } = failingFn(0, rpcErr(LnbitsErrorCode.InternalError), 'ok')
|
||||||
|
expect(await withRetry(fn, { sleep: noSleep })).toBe('ok')
|
||||||
|
expect(calls()).toBe(1)
|
||||||
|
})
|
||||||
|
|
||||||
|
it('retries a transient operator_signer_unavailable, then succeeds', async () => {
|
||||||
|
const { fn, calls } = failingFn(1, rpcErr(LnbitsErrorCode.OperatorSignerUnavailable), 'ok')
|
||||||
|
expect(await withRetry(fn, { sleep: noSleep })).toBe('ok')
|
||||||
|
expect(calls()).toBe(2)
|
||||||
|
})
|
||||||
|
|
||||||
|
it('retries rate_limited (long backoff) up to maxAttempts then throws', async () => {
|
||||||
|
const err = rpcErr(LnbitsErrorCode.RateLimited)
|
||||||
|
const { fn, calls } = failingFn(99, err, 'never')
|
||||||
|
await expect(withRetry(fn, { sleep: noSleep, maxAttempts: 3 })).rejects.toBe(err)
|
||||||
|
expect(calls()).toBe(3)
|
||||||
|
})
|
||||||
|
|
||||||
|
it('internal_error (retry-once) retries exactly once', async () => {
|
||||||
|
const err = rpcErr(LnbitsErrorCode.InternalError)
|
||||||
|
const { fn, calls } = failingFn(99, err, 'never')
|
||||||
|
await expect(withRetry(fn, { sleep: noSleep, maxAttempts: 5 })).rejects.toBe(err)
|
||||||
|
expect(calls()).toBe(2) // initial + one retry, then null delay stops it
|
||||||
|
})
|
||||||
|
|
||||||
|
it('throws a terminal error immediately (no retry)', async () => {
|
||||||
|
const err = rpcErr(LnbitsErrorCode.InsufficientBalance)
|
||||||
|
const { fn, calls } = failingFn(99, err, 'never')
|
||||||
|
await expect(withRetry(fn, { sleep: noSleep })).rejects.toBe(err)
|
||||||
|
expect(calls()).toBe(1)
|
||||||
|
})
|
||||||
|
|
||||||
|
it('treats unauthorized as terminal (no retry)', async () => {
|
||||||
|
const err = rpcErr(LnbitsErrorCode.Unauthorized)
|
||||||
|
const { fn, calls } = failingFn(99, err, 'never')
|
||||||
|
await expect(withRetry(fn, { sleep: noSleep })).rejects.toBe(err)
|
||||||
|
expect(calls()).toBe(1)
|
||||||
|
})
|
||||||
|
|
||||||
|
it('retries a transport timeout error', async () => {
|
||||||
|
const timeout = new Error('LnbitsClient.get_wallet: timeout after 30000ms')
|
||||||
|
const { fn, calls } = failingFn(1, timeout, 'ok')
|
||||||
|
expect(await withRetry(fn, { sleep: noSleep })).toBe('ok')
|
||||||
|
expect(calls()).toBe(2)
|
||||||
|
})
|
||||||
|
|
||||||
|
it('does NOT retry an unknown error (rethrows immediately)', async () => {
|
||||||
|
const boom = new Error('relay socket closed')
|
||||||
|
const { fn, calls } = failingFn(99, boom, 'never')
|
||||||
|
await expect(withRetry(fn, { sleep: noSleep })).rejects.toBe(boom)
|
||||||
|
expect(calls()).toBe(1)
|
||||||
|
})
|
||||||
|
|
||||||
|
it('backs off with increasing delay per attempt (retry-backoff)', async () => {
|
||||||
|
const delays: number[] = []
|
||||||
|
const err = rpcErr(LnbitsErrorCode.OperatorSignerUnavailable)
|
||||||
|
const { fn } = failingFn(99, err, 'never')
|
||||||
|
await expect(
|
||||||
|
withRetry(fn, { sleep: (ms) => { delays.push(ms); return Promise.resolve() }, maxAttempts: 3 })
|
||||||
|
).rejects.toBe(err)
|
||||||
|
expect(delays).toEqual([200, 400]) // before attempt 2 and 3; attempt 3 is last → no 3rd sleep
|
||||||
|
})
|
||||||
|
|
||||||
|
it('invokes onRetry with attempt/delay/error', async () => {
|
||||||
|
const onRetry = vi.fn()
|
||||||
|
const { fn } = failingFn(1, rpcErr(LnbitsErrorCode.OperatorSignerUnavailable), 'ok')
|
||||||
|
await withRetry(fn, { sleep: noSleep, onRetry })
|
||||||
|
expect(onRetry).toHaveBeenCalledOnce()
|
||||||
|
expect(onRetry.mock.calls[0]![0]).toMatchObject({ attempt: 1, delayMs: 200 })
|
||||||
|
})
|
||||||
|
})
|
||||||
|
|
@ -27,6 +27,7 @@ import {
|
||||||
} from '@bitSpire/nostr-client'
|
} from '@bitSpire/nostr-client'
|
||||||
import { verifyEvent } from 'nostr-tools'
|
import { verifyEvent } from 'nostr-tools'
|
||||||
import { LnbitsRpcError } from './error-codes.js'
|
import { LnbitsRpcError } from './error-codes.js'
|
||||||
|
import { withRetry } from './retry.js'
|
||||||
|
|
||||||
import type {
|
import type {
|
||||||
LnbitsConfig,
|
LnbitsConfig,
|
||||||
|
|
@ -155,6 +156,22 @@ export class LnbitsClient {
|
||||||
this.startReplyListener()
|
this.startReplyListener()
|
||||||
}
|
}
|
||||||
|
|
||||||
|
/**
|
||||||
|
* Retry-policy switch for IDEMPOTENT reads only (aiolabs/bitspire#52, Phase D).
|
||||||
|
* Transient failures (operator_signer_unavailable / rate_limited /
|
||||||
|
* internal_error / transport timeout) back off and retry; terminal errors
|
||||||
|
* surface immediately. Never used for pay/create — those would double-pay or
|
||||||
|
* duplicate on retry.
|
||||||
|
*/
|
||||||
|
private idempotent<T>(fn: () => Promise<T>): Promise<T> {
|
||||||
|
return withRetry(fn, {
|
||||||
|
onRetry: ({ attempt, delayMs, error }) => {
|
||||||
|
const code = error instanceof LnbitsRpcError ? error.code : 'timeout'
|
||||||
|
console.warn(`[LnbitsClient] transient ${code} — retry ${attempt} in ${delayMs}ms`)
|
||||||
|
},
|
||||||
|
})
|
||||||
|
}
|
||||||
|
|
||||||
// ============================================================================
|
// ============================================================================
|
||||||
// Wallet
|
// Wallet
|
||||||
// ============================================================================
|
// ============================================================================
|
||||||
|
|
@ -166,8 +183,7 @@ export class LnbitsClient {
|
||||||
*/
|
*/
|
||||||
async getWallet(walletId?: string): Promise<WalletInfo> {
|
async getWallet(walletId?: string): Promise<WalletInfo> {
|
||||||
if (walletId) {
|
if (walletId) {
|
||||||
const data = await this.sendRpc<WalletInfo>('get_wallet', { walletId })
|
return this.idempotent(() => this.sendRpc<WalletInfo>('get_wallet', { walletId }))
|
||||||
return data
|
|
||||||
}
|
}
|
||||||
const wallets = await this.listWallets()
|
const wallets = await this.listWallets()
|
||||||
if (wallets.length === 0) {
|
if (wallets.length === 0) {
|
||||||
|
|
@ -188,7 +204,7 @@ export class LnbitsClient {
|
||||||
|
|
||||||
/** Enumerate every wallet owned by the calling account. */
|
/** Enumerate every wallet owned by the calling account. */
|
||||||
async listWallets(): Promise<WalletInfo[]> {
|
async listWallets(): Promise<WalletInfo[]> {
|
||||||
const data = await this.sendRpc<WalletInfo[]>('list_wallets', {})
|
const data = await this.idempotent(() => this.sendRpc<WalletInfo[]>('list_wallets', {}))
|
||||||
return data ?? []
|
return data ?? []
|
||||||
}
|
}
|
||||||
|
|
||||||
|
|
@ -196,6 +212,10 @@ export class LnbitsClient {
|
||||||
// Invoices
|
// Invoices
|
||||||
// ============================================================================
|
// ============================================================================
|
||||||
|
|
||||||
|
// NOTE: create_invoice / pay_invoice / lnurlw_create_link are NOT wrapped in
|
||||||
|
// `idempotent()` — a retry would mint a duplicate invoice/link or double-pay.
|
||||||
|
// Their errors surface for flow-level handling (state machine / operator).
|
||||||
|
|
||||||
async createInvoice(walletId: string, body: CreateInvoiceBody): Promise<LnbitsPayment> {
|
async createInvoice(walletId: string, body: CreateInvoiceBody): Promise<LnbitsPayment> {
|
||||||
const data = await this.sendRpc<LnbitsPayment>('create_invoice', { walletId, body })
|
const data = await this.sendRpc<LnbitsPayment>('create_invoice', { walletId, body })
|
||||||
return data
|
return data
|
||||||
|
|
@ -208,17 +228,20 @@ export class LnbitsClient {
|
||||||
|
|
||||||
/** Point-lookup of a payment by hash. AUTH_NONE — hashes are hard to guess. */
|
/** Point-lookup of a payment by hash. AUTH_NONE — hashes are hard to guess. */
|
||||||
async getPayment(paymentHash: string): Promise<LnbitsPayment | null> {
|
async getPayment(paymentHash: string): Promise<LnbitsPayment | null> {
|
||||||
const data = await this.sendRpc<LnbitsPayment | null>('get_payment', {
|
const data = await this.idempotent(() =>
|
||||||
body: { payment_hash: paymentHash },
|
this.sendRpc<LnbitsPayment | null>('get_payment', {
|
||||||
})
|
body: { payment_hash: paymentHash },
|
||||||
|
}),
|
||||||
|
)
|
||||||
return data ?? null
|
return data ?? null
|
||||||
}
|
}
|
||||||
|
|
||||||
async decodePayment(paymentRequest: string): Promise<Record<string, unknown>> {
|
async decodePayment(paymentRequest: string): Promise<Record<string, unknown>> {
|
||||||
const data = await this.sendRpc<Record<string, unknown>>('decode_payment', {
|
return this.idempotent(() =>
|
||||||
body: { payment_request: paymentRequest },
|
this.sendRpc<Record<string, unknown>>('decode_payment', {
|
||||||
})
|
body: { payment_request: paymentRequest },
|
||||||
return data
|
}),
|
||||||
|
)
|
||||||
}
|
}
|
||||||
|
|
||||||
// ============================================================================
|
// ============================================================================
|
||||||
|
|
@ -366,18 +389,21 @@ export class LnbitsClient {
|
||||||
}
|
}
|
||||||
|
|
||||||
async getWithdrawLink(walletId: string, id: string): Promise<LnbitsWithdrawLink> {
|
async getWithdrawLink(walletId: string, id: string): Promise<LnbitsWithdrawLink> {
|
||||||
const data = await this.sendRpc<LnbitsWithdrawLink>('lnurlw_get_link', { walletId, body: { id } })
|
return this.idempotent(() =>
|
||||||
return data
|
this.sendRpc<LnbitsWithdrawLink>('lnurlw_get_link', { walletId, body: { id } }),
|
||||||
|
)
|
||||||
}
|
}
|
||||||
|
|
||||||
async listWithdrawLinks(
|
async listWithdrawLinks(
|
||||||
walletId: string | undefined,
|
walletId: string | undefined,
|
||||||
body: { limit?: number; offset?: number } = {},
|
body: { limit?: number; offset?: number } = {},
|
||||||
): Promise<{ data: LnbitsWithdrawLink[]; total: number }> {
|
): Promise<{ data: LnbitsWithdrawLink[]; total: number }> {
|
||||||
return this.sendRpc<{ data: LnbitsWithdrawLink[]; total: number }>('lnurlw_list_links', {
|
return this.idempotent(() =>
|
||||||
walletId,
|
this.sendRpc<{ data: LnbitsWithdrawLink[]; total: number }>('lnurlw_list_links', {
|
||||||
body,
|
walletId,
|
||||||
})
|
body,
|
||||||
|
}),
|
||||||
|
)
|
||||||
}
|
}
|
||||||
|
|
||||||
/**
|
/**
|
||||||
|
|
@ -387,10 +413,12 @@ export class LnbitsClient {
|
||||||
* — the URL a customer wallet GETs to redeem that specific sub-link.
|
* — the URL a customer wallet GETs to redeem that specific sub-link.
|
||||||
*/
|
*/
|
||||||
async getWithdrawLinkUniqueHashes(walletId: string, id: string): Promise<UniqueHashesResponse> {
|
async getWithdrawLinkUniqueHashes(walletId: string, id: string): Promise<UniqueHashesResponse> {
|
||||||
return this.sendRpc<UniqueHashesResponse>('lnurlw_unique_hashes', {
|
return this.idempotent(() =>
|
||||||
walletId,
|
this.sendRpc<UniqueHashesResponse>('lnurlw_unique_hashes', {
|
||||||
body: { id },
|
walletId,
|
||||||
})
|
body: { id },
|
||||||
|
}),
|
||||||
|
)
|
||||||
}
|
}
|
||||||
|
|
||||||
async updateWithdrawLink(
|
async updateWithdrawLink(
|
||||||
|
|
|
||||||
|
|
@ -56,6 +56,8 @@ export {
|
||||||
retryPolicyFor,
|
retryPolicyFor,
|
||||||
} from './error-codes.js'
|
} from './error-codes.js'
|
||||||
export type { RetryPolicy } from './error-codes.js'
|
export type { RetryPolicy } from './error-codes.js'
|
||||||
|
export { withRetry } from './retry.js'
|
||||||
|
export type { WithRetryOptions } from './retry.js'
|
||||||
export type {
|
export type {
|
||||||
LnbitsConfig,
|
LnbitsConfig,
|
||||||
LnbitsRpcRequest,
|
LnbitsRpcRequest,
|
||||||
|
|
|
||||||
79
packages/lnbits/src/retry.ts
Normal file
79
packages/lnbits/src/retry.ts
Normal file
|
|
@ -0,0 +1,79 @@
|
||||||
|
/**
|
||||||
|
* Retry policy switch for the nostr-transport (aiolabs/bitspire#52, Phase D).
|
||||||
|
*
|
||||||
|
* Retries an operation according to the *disposition* of the error it throws —
|
||||||
|
* the machine-readable `retryPolicy` carried by `LnbitsRpcError` (which mirrors
|
||||||
|
* the lnbits canonical enum). Transient conditions back off and retry;
|
||||||
|
* terminal ones throw immediately.
|
||||||
|
*
|
||||||
|
* ⚠️ ONLY wrap IDEMPOTENT operations. A blind retry of `pay_invoice` could
|
||||||
|
* double-pay, and of `create_invoice` / `lnurlw_create_link` would mint
|
||||||
|
* duplicates — those surface their error for flow-level handling (the state
|
||||||
|
* machine / operator) instead. See LnbitsClient for which methods opt in.
|
||||||
|
*/
|
||||||
|
|
||||||
|
import { LnbitsRpcError, type RetryPolicy } from './error-codes.js'
|
||||||
|
|
||||||
|
export interface WithRetryOptions {
|
||||||
|
/** Max total attempts (default 3). */
|
||||||
|
maxAttempts?: number
|
||||||
|
/** Injectable sleep (tests pass a fake-timer-friendly version). */
|
||||||
|
sleep?: (ms: number) => Promise<void>
|
||||||
|
/** Called before each backoff wait — useful for logging. */
|
||||||
|
onRetry?: (info: { attempt: number; delayMs: number; error: unknown }) => void
|
||||||
|
}
|
||||||
|
|
||||||
|
const DEFAULT_MAX_ATTEMPTS = 3
|
||||||
|
|
||||||
|
/** Backoff (ms) for the Nth attempt (1-based), or null if the policy is terminal. */
|
||||||
|
function backoffMs(policy: RetryPolicy, attempt: number): number | null {
|
||||||
|
switch (policy) {
|
||||||
|
case 'retry-backoff':
|
||||||
|
return 200 * 2 ** (attempt - 1) // 200, 400, 800…
|
||||||
|
case 'retry-long-backoff':
|
||||||
|
return 1_000 * 2 ** (attempt - 1) // 1s, 2s, 4s… (rate_limited)
|
||||||
|
case 'retry-once':
|
||||||
|
return attempt === 1 ? 0 : null // exactly one retry (internal_error / absent code)
|
||||||
|
case 'terminal':
|
||||||
|
case 'terminal-idempotent':
|
||||||
|
return null
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
/** Retry delay for an error, or null if it must not be retried. */
|
||||||
|
function delayForError(error: unknown, attempt: number): number | null {
|
||||||
|
if (error instanceof LnbitsRpcError) {
|
||||||
|
return backoffMs(error.retryPolicy, attempt)
|
||||||
|
}
|
||||||
|
// A transport timeout from sendRpc ("…: timeout after <n>ms") is transient.
|
||||||
|
if (error instanceof Error && /timeout after \d+ms/.test(error.message)) {
|
||||||
|
return 200 * 2 ** (attempt - 1)
|
||||||
|
}
|
||||||
|
// Unknown error (programming bug, network teardown) — don't mask it.
|
||||||
|
return null
|
||||||
|
}
|
||||||
|
|
||||||
|
/**
|
||||||
|
* Run `fn`, retrying transient failures per the error's `retryPolicy` (or a
|
||||||
|
* transport timeout) with backoff, up to `maxAttempts`. Terminal errors and
|
||||||
|
* unknown errors throw immediately; the last error is rethrown on exhaustion.
|
||||||
|
*/
|
||||||
|
export async function withRetry<T>(fn: () => Promise<T>, opts: WithRetryOptions = {}): Promise<T> {
|
||||||
|
const maxAttempts = opts.maxAttempts ?? DEFAULT_MAX_ATTEMPTS
|
||||||
|
const sleep = opts.sleep ?? ((ms) => new Promise<void>((resolve) => setTimeout(resolve, ms)))
|
||||||
|
|
||||||
|
let lastError: unknown
|
||||||
|
for (let attempt = 1; attempt <= maxAttempts; attempt++) {
|
||||||
|
try {
|
||||||
|
return await fn()
|
||||||
|
} catch (error) {
|
||||||
|
lastError = error
|
||||||
|
if (attempt === maxAttempts) break
|
||||||
|
const delayMs = delayForError(error, attempt)
|
||||||
|
if (delayMs === null) throw error
|
||||||
|
opts.onRetry?.({ attempt, delayMs, error })
|
||||||
|
await sleep(delayMs)
|
||||||
|
}
|
||||||
|
}
|
||||||
|
throw lastError
|
||||||
|
}
|
||||||
Loading…
Add table
Add a link
Reference in a new issue