feat(machine): durable dispense-report outbox to spirekeeper (ADR-005 §2)
Every cash-out now produces one report_dispense — on success as well as failure — and the machine does not stop sending it until spirekeeper acknowledges it. state.db gains a dispense_reports table (migration v13 → v14): the report is written INSIDE recordTransaction's SQLite transaction, alongside the transactions row, so a crash between the two cannot lose it. Rows carry attempts / last_attempt_at / last_error / acked_at. Three IPC calls (pending / ack / note-attempt) expose it to the renderer. The store builds the report when a cash-out reaches complete, dispenseFault or outOfCash: txid, payment hash, dispense_confirmed, error / error_code / raw_code / error_class, per-denomination requested vs dispensed vs rejected, the per-bay cassette record verbatim, and counts_uncertain. The success report is what lets the server capture (distribute) the settlement; the failure report is what puts a customer on the owed-cash worklist instead of leaving the only record on the ATM. Delivery is at-least-once: a flusher drains pending rows after each persist, on relay (re)connect, and every 60 s, acking only on an OK reply and backing off 30 s · 2^attempts (capped 1 h) otherwise. While spirekeeper has not registered the RPC every send fails the same way; the backoff keeps that quiet and the rows wait — this half ships first. The lightning service exposes reportDispense; the function pointer is set at all three lightning-init sites so the flusher works on every path.
This commit is contained in:
parent
7120f306b6
commit
e8106b665a
7 changed files with 319 additions and 2 deletions
|
|
@ -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<DispenseReportAck>
|
||||
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
|
||||
},
|
||||
|
|
|
|||
|
|
@ -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<Record<number, number> | 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<number, number>()
|
||||
for (const a of ctx.dispenseAmounts) {
|
||||
requestedByDenom.set(a.denomination, (requestedByDenom.get(a.denomination) ?? 0) + a.count)
|
||||
}
|
||||
const seen = new Set<number>()
|
||||
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<void> {
|
||||
if (isElectron && window.electronAPI) {
|
||||
try {
|
||||
|
|
@ -401,6 +459,12 @@ export const useAtmStore = defineStore('atm', () => {
|
|||
)
|
||||
let stopBalanceWatch: (() => void) | null = null
|
||||
let pricePollingInterval: ReturnType<typeof setInterval> | 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<unknown>) | null = null
|
||||
let dispenseReportFlushInterval: ReturnType<typeof setInterval> | 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<void> {
|
||||
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
|
||||
|
|
|
|||
13
apps/machine/src/types/electron.d.ts
vendored
13
apps/machine/src/types/electron.d.ts
vendored
|
|
@ -164,6 +164,19 @@ declare global {
|
|||
since: number
|
||||
}) => Promise<{ reason: string; errorCode: string | null; rawCode: string | null; since: number }>
|
||||
clearCashOutHold: () => Promise<boolean>
|
||||
// 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<boolean>
|
||||
noteDispenseReportAttempt: (txid: string, error: string | null) => Promise<void>
|
||||
markStatePublished: (unixTimestamp: number) => Promise<void>
|
||||
saveBunkerBinding: (binding: BunkerBindingRecord) => Promise<void>
|
||||
clearBunkerBinding: () => Promise<void>
|
||||
|
|
|
|||
|
|
@ -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 {
|
||||
|
|
|
|||
Loading…
Add table
Add a link
Reference in a new issue