diff --git a/apps/machine/src/services/operator-config.ts b/apps/machine/src/services/operator-config.ts index ba9ddcb..d853114 100644 --- a/apps/machine/src/services/operator-config.ts +++ b/apps/machine/src/services/operator-config.ts @@ -56,6 +56,17 @@ export interface OperatorConfigServiceConfig { export interface OperatorConfigService { /** Unsubscribe from operator events and free resources. */ 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 +} + +const NOOP_SERVICE: OperatorConfigService = { + stop: () => {}, + publishCassettesState: async () => {}, } export async function startOperatorConfigService( @@ -63,11 +74,11 @@ export async function startOperatorConfigService( ): Promise { if (cfg.operatorPubkeys.length === 0) { console.log('[OperatorConfig] No operator pubkeys configured — service disabled') - return { stop: () => {} } + return NOOP_SERVICE } if (!isElectron || !window.electronAPI) { console.log('[OperatorConfig] Not in Electron — service disabled (browser dev mode)') - return { stop: () => {} } + return NOOP_SERVICE } const api = window.electronAPI const machineId = cfg.machineId ?? cfg.signer.pubkey @@ -103,6 +114,12 @@ export async function startOperatorConfigService( return { 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( `[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:`), 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, api: NonNullable, machineId: string -): Promise { - const already = await api.getBootstrapPublishedAt() - if (already !== null) { - console.log('[OperatorConfig] Bootstrap already published at unix', already) - return - } +): Promise { const cassettes = await api.loadCassettes() - if (cassettes.length === 0) { - console.log('[OperatorConfig] state.db.cassettes empty — skipping bootstrap') - return - } - + if (cassettes.length === 0) return false const operatorPubkey = cfg.operatorPubkeys[0] - if (!operatorPubkey) { - console.log('[OperatorConfig] No operator pubkey — skipping bootstrap') - return - } + if (!operatorPubkey) return false const positions: Record = {} for (const c of cassettes) { @@ -235,6 +258,30 @@ async function maybePublishBootstrap( }) await cfg.nostrClient.publish(event) - await api.markBootstrapPublished(Math.floor(Date.now() / 1000)) - console.log('[OperatorConfig] Bootstrap hello-event published:', { dTag, eventId: event.id }) + console.log('[OperatorConfig] cassettes-state 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, + machineId: string +): Promise { + 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') + } } diff --git a/apps/machine/src/stores/atm.ts b/apps/machine/src/stores/atm.ts index ca22847..2bfaeea 100644 --- a/apps/machine/src/stores/atm.ts +++ b/apps/machine/src/stores/atm.ts @@ -521,7 +521,10 @@ export const useAtmStore = defineStore('atm', () => { bills, cassettes: dr?.cassettes, 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, cassettes: dr?.cassettes, 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())) } } diff --git a/packages/lnbits/src/__tests__/retry.test.ts b/packages/lnbits/src/__tests__/retry.test.ts new file mode 100644 index 0000000..8a57b45 --- /dev/null +++ b/packages/lnbits/src/__tests__/retry.test.ts @@ -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(failTimes: number, err: unknown, value: T): { fn: () => Promise; 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 }) + }) +}) diff --git a/packages/lnbits/src/client.ts b/packages/lnbits/src/client.ts index 946df5d..de408a1 100644 --- a/packages/lnbits/src/client.ts +++ b/packages/lnbits/src/client.ts @@ -27,6 +27,7 @@ import { } from '@bitSpire/nostr-client' import { verifyEvent } from 'nostr-tools' import { LnbitsRpcError } from './error-codes.js' +import { withRetry } from './retry.js' import type { LnbitsConfig, @@ -155,6 +156,22 @@ export class LnbitsClient { 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(fn: () => Promise): Promise { + 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 // ============================================================================ @@ -166,8 +183,7 @@ export class LnbitsClient { */ async getWallet(walletId?: string): Promise { if (walletId) { - const data = await this.sendRpc('get_wallet', { walletId }) - return data + return this.idempotent(() => this.sendRpc('get_wallet', { walletId })) } const wallets = await this.listWallets() if (wallets.length === 0) { @@ -188,7 +204,7 @@ export class LnbitsClient { /** Enumerate every wallet owned by the calling account. */ async listWallets(): Promise { - const data = await this.sendRpc('list_wallets', {}) + const data = await this.idempotent(() => this.sendRpc('list_wallets', {})) return data ?? [] } @@ -196,6 +212,10 @@ export class LnbitsClient { // 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 { const data = await this.sendRpc('create_invoice', { walletId, body }) return data @@ -208,17 +228,20 @@ export class LnbitsClient { /** Point-lookup of a payment by hash. AUTH_NONE — hashes are hard to guess. */ async getPayment(paymentHash: string): Promise { - const data = await this.sendRpc('get_payment', { - body: { payment_hash: paymentHash }, - }) + const data = await this.idempotent(() => + this.sendRpc('get_payment', { + body: { payment_hash: paymentHash }, + }), + ) return data ?? null } async decodePayment(paymentRequest: string): Promise> { - const data = await this.sendRpc>('decode_payment', { - body: { payment_request: paymentRequest }, - }) - return data + return this.idempotent(() => + this.sendRpc>('decode_payment', { + body: { payment_request: paymentRequest }, + }), + ) } // ============================================================================ @@ -366,18 +389,21 @@ export class LnbitsClient { } async getWithdrawLink(walletId: string, id: string): Promise { - const data = await this.sendRpc('lnurlw_get_link', { walletId, body: { id } }) - return data + return this.idempotent(() => + this.sendRpc('lnurlw_get_link', { walletId, body: { id } }), + ) } async listWithdrawLinks( walletId: string | undefined, body: { limit?: number; offset?: number } = {}, ): Promise<{ data: LnbitsWithdrawLink[]; total: number }> { - return this.sendRpc<{ data: LnbitsWithdrawLink[]; total: number }>('lnurlw_list_links', { - walletId, - body, - }) + return this.idempotent(() => + this.sendRpc<{ data: LnbitsWithdrawLink[]; total: number }>('lnurlw_list_links', { + walletId, + body, + }), + ) } /** @@ -387,10 +413,12 @@ export class LnbitsClient { * — the URL a customer wallet GETs to redeem that specific sub-link. */ async getWithdrawLinkUniqueHashes(walletId: string, id: string): Promise { - return this.sendRpc('lnurlw_unique_hashes', { - walletId, - body: { id }, - }) + return this.idempotent(() => + this.sendRpc('lnurlw_unique_hashes', { + walletId, + body: { id }, + }), + ) } async updateWithdrawLink( diff --git a/packages/lnbits/src/index.ts b/packages/lnbits/src/index.ts index 10e0f05..2085593 100644 --- a/packages/lnbits/src/index.ts +++ b/packages/lnbits/src/index.ts @@ -56,6 +56,8 @@ export { retryPolicyFor, } from './error-codes.js' export type { RetryPolicy } from './error-codes.js' +export { withRetry } from './retry.js' +export type { WithRetryOptions } from './retry.js' export type { LnbitsConfig, LnbitsRpcRequest, diff --git a/packages/lnbits/src/retry.ts b/packages/lnbits/src/retry.ts new file mode 100644 index 0000000..0495859 --- /dev/null +++ b/packages/lnbits/src/retry.ts @@ -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 + /** 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 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(fn: () => Promise, opts: WithRetryOptions = {}): Promise { + const maxAttempts = opts.maxAttempts ?? DEFAULT_MAX_ATTEMPTS + const sleep = opts.sleep ?? ((ms) => new Promise((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 +}