diff --git a/apps/machine/electron/main.ts b/apps/machine/electron/main.ts index 0295894..68b9f88 100644 --- a/apps/machine/electron/main.ts +++ b/apps/machine/electron/main.ts @@ -31,6 +31,9 @@ import { setCashOutHold, clearCashOutHold, type CashOutHold, + pendingDispenseReports, + markDispenseReportAcked, + noteDispenseReportAttempt, markStatePublished, resetStatePublishWatermark, resetForRepair, @@ -578,6 +581,22 @@ ipcMain.handle('state:set-cash-out-hold', (_event, hold: CashOutHold): CashOutHo return setCashOutHold(hold) }) ipcMain.handle('state:clear-cash-out-hold', (): boolean => clearCashOutHold()) + +// Dispense-report outbox (ADR-005 §2) — at-least-once to spirekeeper +ipcMain.handle('state:pending-dispense-reports', (_event, limit?: number) => + pendingDispenseReports(typeof limit === 'number' ? limit : 20) +) +ipcMain.handle('state:ack-dispense-report', (_event, txid: string): boolean => { + if (typeof txid !== 'string' || !txid) throw new Error('Invalid txid') + return markDispenseReportAcked(txid) +}) +ipcMain.handle( + 'state:note-dispense-report-attempt', + (_event, txid: string, error: string | null): void => { + if (typeof txid !== 'string' || !txid) throw new Error('Invalid txid') + noteDispenseReportAttempt(txid, typeof error === 'string' ? error.slice(0, 512) : null) + } +) ipcMain.handle('state:mark-state-published', (_event, unixTimestamp: number): void => { markStatePublished(unixTimestamp) }) diff --git a/apps/machine/electron/preload.ts b/apps/machine/electron/preload.ts index 83ab409..bc75046 100644 --- a/apps/machine/electron/preload.ts +++ b/apps/machine/electron/preload.ts @@ -15,6 +15,16 @@ interface CashOutHold { since: number } +/** Mirrors state-store.PendingDispenseReport (ADR-005 §2). */ +interface PendingDispenseReport { + txid: string + payload: unknown + createdAt: number + attempts: number + lastAttemptAt: number | null + lastError: string | null +} + /** * Runtime configuration interface (public info only) * These values are read from environment variables at runtime (not build time) @@ -127,6 +137,13 @@ contextBridge.exposeInMainWorld('electronAPI', { setCashOutHold: (hold: CashOutHold): Promise => ipcRenderer.invoke('state:set-cash-out-hold', hold), clearCashOutHold: (): Promise => ipcRenderer.invoke('state:clear-cash-out-hold'), + // Dispense-report outbox (ADR-005 §2) + pendingDispenseReports: (limit?: number): Promise => + ipcRenderer.invoke('state:pending-dispense-reports', limit), + ackDispenseReport: (txid: string): Promise => + ipcRenderer.invoke('state:ack-dispense-report', txid), + noteDispenseReportAttempt: (txid: string, error: string | null): Promise => + ipcRenderer.invoke('state:note-dispense-report-attempt', txid, error), markStatePublished: (unixTimestamp: number): Promise => ipcRenderer.invoke('state:mark-state-published', unixTimestamp), @@ -323,6 +340,9 @@ declare global { getCashOutHold: () => Promise setCashOutHold: (hold: CashOutHold) => Promise clearCashOutHold: () => Promise + pendingDispenseReports: (limit?: number) => Promise + ackDispenseReport: (txid: string) => Promise + noteDispenseReportAttempt: (txid: string, error: string | null) => Promise markStatePublished: (unixTimestamp: number) => Promise saveBunkerBinding: (binding: BunkerBindingRecord) => Promise clearBunkerBinding: () => Promise diff --git a/apps/machine/electron/state-store.ts b/apps/machine/electron/state-store.ts index ec4d682..67afd49 100644 --- a/apps/machine/electron/state-store.ts +++ b/apps/machine/electron/state-store.ts @@ -10,12 +10,13 @@ */ import Database from 'better-sqlite3' +import type { DispenseReportBody } from '@bitSpire/lnbits' import path from 'node:path' import fs from 'node:fs' let db: Database.Database | null = null -const SCHEMA_VERSION = '13' +const SCHEMA_VERSION = '14' function getDbPath(): string { const prodDir = '/var/lib/bitspire' @@ -136,6 +137,16 @@ export function initDatabase(dbPath?: string): void { relays TEXT, lnbits_server_pubkey TEXT ); + + CREATE TABLE IF NOT EXISTS dispense_reports ( + txid TEXT PRIMARY KEY REFERENCES transactions(txid), + payload TEXT NOT NULL, + created_at INTEGER NOT NULL, + attempts INTEGER NOT NULL DEFAULT 0, + last_attempt_at INTEGER, + last_error TEXT, + acked_at INTEGER + ); `) // Seed meta + cashbox if first run, or run migrations @@ -418,6 +429,32 @@ export function initDatabase(dbPath?: string): void { existing.value = '13' } + if (existing && existing.value === '13') { + // Migration v13 → v14: the dispense-report outbox (ADR-005 §2). + // + // Every cash-out's outcome — success or failure — is reported to + // spirekeeper over a kind-21000 RPC, and that report is what lets the + // server capture (distribute) the settlement or surface a customer who + // is owed cash. A relay gives the publisher no delivery guarantee, so the + // report is written here, in the SAME transaction as the transactions + // row, and resent until the server acknowledges it. Idempotent on txid + // server-side; `attempts` / `last_error` drive the resend backoff. + db.exec(` + CREATE TABLE IF NOT EXISTS dispense_reports ( + txid TEXT PRIMARY KEY REFERENCES transactions(txid), + payload TEXT NOT NULL, + created_at INTEGER NOT NULL, + attempts INTEGER NOT NULL DEFAULT 0, + last_attempt_at INTEGER, + last_error TEXT, + acked_at INTEGER + ); + `) + db.prepare('UPDATE meta SET value = ? WHERE key = ?').run('14', 'schema_version') + console.log('[StateStore] Migrated schema v13 → v14 (added dispense_reports outbox)') + existing.value = '14' + } + // Defensive: a fresh install at SCHEMA_VERSION skips all migrations. // Seed the operator-config meta rows if they're missing (idempotent). const seedMeta = db.prepare('INSERT OR IGNORE INTO meta (key, value) VALUES (?, ?)') @@ -580,6 +617,70 @@ export function clearCashOutHold(): boolean { return had } +// --------------------------------------------------------------------------- +// Dispense-report outbox (ADR-005 §2) +// --------------------------------------------------------------------------- + +export interface PendingDispenseReport { + txid: string + payload: DispenseReportBody + createdAt: number + attempts: number + lastAttemptAt: number | null + lastError: string | null +} + +/** Unacknowledged reports, oldest first. The renderer applies the backoff. */ +export function pendingDispenseReports(limit = 20): PendingDispenseReport[] { + if (!db) throw new Error('Database not initialized') + const rows = db + .prepare( + 'SELECT txid, payload, created_at, attempts, last_attempt_at, last_error FROM dispense_reports WHERE acked_at IS NULL ORDER BY created_at ASC LIMIT ?' + ) + .all(limit) as Array<{ + txid: string + payload: string + created_at: number + attempts: number + last_attempt_at: number | null + last_error: string | null + }> + const out: PendingDispenseReport[] = [] + for (const r of rows) { + try { + out.push({ + txid: r.txid, + payload: JSON.parse(r.payload) as DispenseReportBody, + createdAt: r.created_at, + attempts: r.attempts, + lastAttemptAt: r.last_attempt_at, + lastError: r.last_error, + }) + } catch { + console.error('[StateStore] dispense_reports row has unparseable payload:', r.txid) + } + } + return out +} + +/** The server acknowledged this report. Returns whether a row changed. */ +export function markDispenseReportAcked(txid: string): boolean { + if (!db) throw new Error('Database not initialized') + const res = db + .prepare('UPDATE dispense_reports SET acked_at = ? WHERE txid = ? AND acked_at IS NULL') + .run(Date.now(), txid) + if (res.changes > 0) console.log('[StateStore] Dispense report acked:', txid) + return res.changes > 0 +} + +/** A send was attempted and did not get an OK. Drives the resend backoff. */ +export function noteDispenseReportAttempt(txid: string, error: string | null): void { + if (!db) throw new Error('Database not initialized') + db.prepare( + 'UPDATE dispense_reports SET attempts = attempts + 1, last_attempt_at = ?, last_error = ? WHERE txid = ?' + ).run(Date.now(), error, txid) +} + /** * A counter bumped on every local change to a bay count, from any cause. * @@ -1241,6 +1342,11 @@ interface TransactionInput { rejected: number }[] error?: string | null + /** + * ADR-005 §2: dispense outcome to queue for spirekeeper. Inserted in the + * same transaction as the row so a crash between them cannot lose it. + */ + report?: DispenseReportBody } /** @@ -1262,6 +1368,12 @@ export function recordTransaction(tx: TransactionInput): void { const insertBill = db.prepare( 'INSERT INTO transaction_bills (txid, denomination, count) VALUES (?, ?, ?)' ) + // Outbox row (ADR-005 §2). REPLACE: a re-record of the same txid (should not + // happen, but a crash-replay could) refreshes the payload and resets the + // delivery state rather than failing the whole transaction. + const insertReport = db.prepare( + 'INSERT OR REPLACE INTO dispense_reports (txid, payload, created_at, attempts, last_attempt_at, last_error, acked_at) VALUES (?, ?, ?, 0, NULL, NULL, NULL)' + ) const insertCassetteBill = db.prepare( 'INSERT INTO cassette_bills (txid, name, position, denomination, provisioned, dispensed, rejected) VALUES (?, ?, ?, ?, ?, ?, ?)' ) @@ -1294,6 +1406,10 @@ export function recordTransaction(tx: TransactionInput): void { insertBill.run(t.txid, bill.denomination, bill.count) } + if (t.report) { + insertReport.run(t.txid, JSON.stringify(t.report), Date.now()) + } + // Insert per-cassette detail when available if (t.cassettes) { for (const c of t.cassettes) { diff --git a/apps/machine/src/services/lightning.ts b/apps/machine/src/services/lightning.ts index 3a72f63..10d2566 100644 --- a/apps/machine/src/services/lightning.ts +++ b/apps/machine/src/services/lightning.ts @@ -14,7 +14,7 @@ import { NostrClient, type Signer } from '@bitSpire/nostr-client' import { resolveSigner } from './signer-resolver.js' -import { LnbitsClient } from '@bitSpire/lnbits' +import { LnbitsClient, type DispenseReportBody, type DispenseReportAck } from '@bitSpire/lnbits' import { CLINKClient } from '@bitSpire/clink' import type { OfferRequest, ManagementRequest, ManagementResponse } from '@bitSpire/clink' import type { ATMServices, ATMContext } from '@bitSpire/state-machine' @@ -223,6 +223,8 @@ export interface LightningBackend { } interface LightningServices { + /** ADR-005 §2: send one cash-out's dispense outcome to spirekeeper (outbox-driven). */ + reportDispense: (body: DispenseReportBody) => Promise nostrClient: NostrClient lightningPub: LightningBackend clink: CLINKClient @@ -686,6 +688,11 @@ export async function initializeLightningServices(options?: { signer, operatorPubkeys: CONFIG.operatorPubkeys, atmServices, + /** + * ADR-005 §2: send one cash-out's dispense outcome to spirekeeper. The + * store keeps these in a durable outbox and calls this until it resolves. + */ + reportDispense: (body: DispenseReportBody) => lnbits.reportDispense(body), onOfferRequest: (callback: OfferRequestCallback) => { offerRequestCallback = callback }, diff --git a/apps/machine/src/stores/atm.ts b/apps/machine/src/stores/atm.ts index 8131c85..4ed884f 100644 --- a/apps/machine/src/stores/atm.ts +++ b/apps/machine/src/stores/atm.ts @@ -1,4 +1,5 @@ import { defineStore } from 'pinia' +import type { DispenseReportBody } from '@bitSpire/lnbits' import { ref, computed, watch } from 'vue' import { createATMMachine, @@ -9,6 +10,7 @@ import { type SnapshotFrom, type ATMMachine, type AccessRole, + type DispenseCashResult, } from '@bitSpire/state-machine' import type { AccessControlConfig, CardSession } from '@/types/electron' import { initializeLightningServices, fetchBtcPrice } from '@/services/lightning' @@ -192,6 +194,62 @@ async function loadInventoryFromDb(): Promise | null> { /** * Persist a completed transaction to SQLite via IPC. */ +/** + * ADR-005 §2: the dispense outcome spirekeeper captures on. Built from the + * machine context at the moment the terminal state is entered; written in the + * same SQLite transaction as the transactions row (see TransactionRecord.report). + * `requested` per denomination comes from the sale, `dispensed`/`rejected` from + * the hardware report; `cassettes` is the per-bay record verbatim. + */ +function buildDispenseReport( + ctx: ATMContext, + dr: DispenseCashResult | null, + countsUncertain: boolean +): DispenseReportBody { + const requestedByDenom = new Map() + for (const a of ctx.dispenseAmounts) { + requestedByDenom.set(a.denomination, (requestedByDenom.get(a.denomination) ?? 0) + a.count) + } + const seen = new Set() + const bills: DispenseReportBody['bills'] = [] + for (const b of dr?.bills ?? []) { + seen.add(b.denomination) + bills.push({ + denomination: b.denomination, + requested: requestedByDenom.get(b.denomination) ?? 0, + dispensed: b.dispensed, + rejected: b.rejected, + }) + } + // A denomination that was asked for but never appears in the report (the + // dispenser threw before reporting) still needs a row: requested, zero out. + for (const [denomination, requested] of requestedByDenom) { + if (!seen.has(denomination)) bills.push({ denomination, requested, dispensed: 0, rejected: 0 }) + } + return { + txid: ctx.txid ?? '', + payment_hash: ctx.paymentHash, + tx_type: 'cash_out', + dispense_confirmed: dr?.dispenseConfirmed === true, + error: dr?.error ?? ctx.error ?? null, + error_code: dr?.errorCode ?? (ctx.error && !dr ? 'DispenseThrew' : null), + raw_code: dr?.rawCode ?? null, + error_class: dr?.errorClass ?? (ctx.error && !dr ? 'terminal' : null), + fiat_cents: ctx.fiatCents, + currency: ctx.currency, + bills, + cassettes: (dr?.cassettes ?? []).map((c) => ({ + position: c.position, + denomination: c.denomination, + provisioned: c.provisioned, + dispensed: c.dispensed, + rejected: c.rejected, + })), + counts_uncertain: countsUncertain, + at: Math.floor(Date.now() / 1000), + } +} + async function persistTransaction(tx: TransactionRecord): Promise { if (isElectron && window.electronAPI) { try { @@ -401,6 +459,12 @@ export const useAtmStore = defineStore('atm', () => { ) let stopBalanceWatch: (() => void) | null = null let pricePollingInterval: ReturnType | null = null + // ADR-005 §2 — dispense-report outbox. The function pointer is set at every + // lightning-init site; the flusher drains state.db's dispense_reports table + // to spirekeeper with backoff until each row is acked. + let reportDispenseFn: ((body: DispenseReportBody) => Promise) | null = null + let dispenseReportFlushInterval: ReturnType | null = null + let dispenseReportFlushing = false // Store reference to ATM services for direct calls let atmServicesRef: ATMServices | null = null @@ -684,6 +748,11 @@ export const useAtmStore = defineStore('atm', () => { .catch((e) => console.warn('[ATM] Could not flag counts unverified:', e)) } + const countsUncertain = + !dr || + (dr.bills.reduce((sum, b) => sum + b.dispensed, 0) === 0 && + !!dr.error && + dr.errorClass !== 'inventory') persistTransaction({ txid: ctx.txid, type: 'cash_out', @@ -697,11 +766,14 @@ export const useAtmStore = defineStore('atm', () => { bills, cassettes: dr?.cassettes, error: dr?.error ?? ctx.error, + // ADR-005 §2 — queued in the same SQLite transaction as the row. + report: buildDispenseReport(ctx, dr, countsUncertain), }) .then(() => reloadPersistedInventory()) // Republish cassette state — a partial dispense changed counts, and // the payload now carries the hold / unverified flags. .then(() => operatorConfigSvc?.publishCassettesState()) + .then(() => flushDispenseReports()) } } @@ -731,11 +803,15 @@ export const useAtmStore = defineStore('atm', () => { bills, cassettes: dr?.cassettes, error: dr?.error, + // ADR-005 §2 — the SUCCESS report is what lets spirekeeper capture + // (distribute) the settlement. Cash-out only; cash-in has no dispense. + ...(isCashInTx ? {} : { report: buildDispenseReport(ctx, dr, false) }), }) .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())) + .then(() => (isCashInTx ? undefined : flushDispenseReports())) } } @@ -1084,6 +1160,8 @@ export const useAtmStore = defineStore('atm', () => { try { const services = await initializeLightningServices({ strict: !allowMockFallback.value }) + reportDispenseFn = services.reportDispense + startDispenseReportFlusher() useLiveServices.value = true connectionStatus.value = 'connected' console.log('[ATM] Connected to Lightning.Pub!') @@ -1096,6 +1174,8 @@ export const useAtmStore = defineStore('atm', () => { services.nostrClient.on('connect', () => { connectionStatus.value = 'connected' console.log('[ATM] Relay reconnected') + // A report queued during the outage goes now, not at the next tick. + void flushDispenseReports() }) // Store references to clients for direct operations @@ -1383,6 +1463,8 @@ export const useAtmStore = defineStore('atm', () => { // Initialize Lightning services const lightning = await initializeLightningServices({ strict: !allowMockFallback.value }) + reportDispenseFn = lightning.reportDispense + startDispenseReportFlusher() useLiveServices.value = true connectionStatus.value = 'connected' lightningPub.value = lightning.lightningPub @@ -1676,6 +1758,8 @@ export const useAtmStore = defineStore('atm', () => { // Initialize Lightning services const lightning = await initializeLightningServices({ strict: !allowMockFallback.value }) + reportDispenseFn = lightning.reportDispense + startDispenseReportFlusher() useLiveServices.value = true connectionStatus.value = 'connected' lightningPub.value = lightning.lightningPub @@ -2044,7 +2128,59 @@ export const useAtmStore = defineStore('atm', () => { pricePollingInterval = setInterval(poll, 30_000) } + /** + * Drain the dispense-report outbox (ADR-005 §2). At-least-once: a row is + * acked only on an OK reply; anything else bumps `attempts` and the row is + * retried after an exponential backoff (30 s · 2^attempts, capped at 1 h). + * While spirekeeper has not registered `report_dispense` every send fails + * the same way — the backoff keeps that from being noisy, and the rows wait. + * Triggers: right after each persist, on relay (re)connect, every 60 s. + */ + async function flushDispenseReports(): Promise { + if (!isElectron || !window.electronAPI || !reportDispenseFn) return + if (dispenseReportFlushing) return + dispenseReportFlushing = true + try { + const pending = await window.electronAPI.pendingDispenseReports(20) + const now = Date.now() + for (const row of pending) { + const backoffMs = Math.min(30_000 * 2 ** row.attempts, 3_600_000) + if (row.lastAttemptAt && row.lastAttemptAt + backoffMs > now) continue + try { + await reportDispenseFn(row.payload as DispenseReportBody) + await window.electronAPI.ackDispenseReport(row.txid) + console.log(`[ATM] Dispense report delivered: ${row.txid}`) + } catch (e) { + const msg = e instanceof Error ? e.message : String(e) + await window.electronAPI.noteDispenseReportAttempt(row.txid, msg) + console.warn( + `[ATM] Dispense report ${row.txid} not delivered (attempt ${row.attempts + 1}): ${msg}` + ) + } + } + } catch (e) { + console.warn('[ATM] Dispense-report flush failed:', e) + } finally { + dispenseReportFlushing = false + } + } + + function startDispenseReportFlusher() { + if (dispenseReportFlushInterval) return + dispenseReportFlushInterval = setInterval(() => void flushDispenseReports(), 60_000) + void flushDispenseReports() + } + + function stopDispenseReportFlusher() { + if (dispenseReportFlushInterval) { + clearInterval(dispenseReportFlushInterval) + dispenseReportFlushInterval = null + } + } + function stopPricePolling() { + // Both are store-lifetime intervals; whoever stops one stops the other. + stopDispenseReportFlusher() if (pricePollingInterval) { clearInterval(pricePollingInterval) pricePollingInterval = null diff --git a/apps/machine/src/types/electron.d.ts b/apps/machine/src/types/electron.d.ts index 9c54e75..f3b4db5 100644 --- a/apps/machine/src/types/electron.d.ts +++ b/apps/machine/src/types/electron.d.ts @@ -164,6 +164,19 @@ declare global { since: number }) => Promise<{ reason: string; errorCode: string | null; rawCode: string | null; since: number }> clearCashOutHold: () => Promise + // Dispense-report outbox (ADR-005 §2) + pendingDispenseReports: (limit?: number) => Promise< + Array<{ + txid: string + payload: unknown + createdAt: number + attempts: number + lastAttemptAt: number | null + lastError: string | null + }> + > + ackDispenseReport: (txid: string) => Promise + noteDispenseReportAttempt: (txid: string, error: string | null) => Promise markStatePublished: (unixTimestamp: number) => Promise saveBunkerBinding: (binding: BunkerBindingRecord) => Promise clearBunkerBinding: () => Promise diff --git a/apps/machine/src/types/state.ts b/apps/machine/src/types/state.ts index 7d0d9c3..9b780c3 100644 --- a/apps/machine/src/types/state.ts +++ b/apps/machine/src/types/state.ts @@ -34,6 +34,12 @@ export interface TransactionRecord { }[] error?: string | null remediatedBy?: string | null + /** + * ADR-005 §2: the dispense outcome to queue for spirekeeper, written in + * the same SQLite transaction as the row so a crash between the two cannot + * lose it. Cash-out only. Shape is @bitSpire/lnbits DispenseReportBody. + */ + report?: import('@bitSpire/lnbits').DispenseReportBody } export interface ATMAvailability {