diff --git a/apps/machine/electron/__tests__/state-store-transactions.test.ts b/apps/machine/electron/__tests__/state-store-transactions.test.ts index 93c3434..02dd87e 100644 --- a/apps/machine/electron/__tests__/state-store-transactions.test.ts +++ b/apps/machine/electron/__tests__/state-store-transactions.test.ts @@ -12,8 +12,10 @@ import { afterEach, beforeEach, describe, expect, it } from 'vitest' import { - applyOperatorCassettesConfig, + applyOperatorCassetteOps, closeDatabase, + getAppliedOpIds, + getCassetteStateSeq, getCashbox, getCountsUncertainSince, getInventory, @@ -337,31 +339,157 @@ describe('state-store: unverified counts after a silent dispense', () => { expect(getCountsUncertainSince()).toBe(1000) }) - it('clears when an operator asserts real counts', () => { + it('clears on a recount, because that is what a recount is', () => { markCountsUncertain(1000) - const applied = applyOperatorCassettesConfig( - { - positions: { - '1': { denomination: 20, count: 40 }, - '2': { denomination: 20, count: 40 }, - '3': { denomination: 50, count: 25 }, - }, - }, - 1_700_000_000 - ) - expect(applied.applied).toBe(true) - // A recount is exactly the operator asserting authoritative counts. + const result = applyOperatorCassetteOps([ + { id: 'op-1', at: 1_700_000_000, type: 'recount', position: 1, count: 40 }, + ]) + expect(result.applied).toEqual(['op-1']) expect(getCountsUncertainSince()).toBeNull() }) - it('leaves the flag alone when the operator config is rejected', () => { + it('does not clear on a refill', () => { + // A refill adds to a number still known to be wrong. Only someone + // opening the bay and counting it resolves that. markCountsUncertain(1000) - // Position key-set mismatch — the bay layout is hardware-determined. - const applied = applyOperatorCassettesConfig( - { positions: { '1': { denomination: 20, count: 40 } } }, - 1_700_000_001 - ) - expect(applied.applied).toBe(false) + const result = applyOperatorCassetteOps([ + { id: 'op-1', at: 1_700_000_000, type: 'refill', position: 1, bills: 10 }, + ]) + expect(result.applied).toEqual(['op-1']) + expect(getCountsUncertainSince()).toBe(1000) + }) + + it('leaves the flag alone when the op is rejected', () => { + markCountsUncertain(1000) + // Bay 9 does not exist — the layout is hardware-determined. + const result = applyOperatorCassetteOps([ + { id: 'op-1', at: 1_700_000_000, type: 'recount', position: 9, count: 40 }, + ]) + expect(result.applied).toEqual([]) + expect(result.rejected).toHaveLength(1) expect(getCountsUncertainSince()).toBe(1000) }) }) + +describe('state-store: operator cassette operations (ADR-004)', () => { + beforeEach(() => { + seedDuplicateDenomBays() + }) + + it('applies a refill as a delta, not a total', () => { + applyOperatorCassetteOps([ + { id: 'op-1', at: 1_700_000_000, type: 'refill', position: 1, bills: 30 }, + ]) + expect(loadCassettes().find((c) => c.position === 1)!.count).toBe(80) + }) + + it('is a no-op on a re-delivered operation', () => { + // Addressable events are re-delivered on every relay reconnect and the + // operator republishes a WINDOW, so the same op arrives many times. A + // delta applied twice is simply wrong, which is why every op carries an + // id and this table records the ones already applied. + const op = { + id: 'op-1', + at: 1_700_000_000, + type: 'refill' as const, + position: 1, + bills: 30, + } + expect(applyOperatorCassetteOps([op]).applied).toEqual(['op-1']) + expect(applyOperatorCassetteOps([op]).applied).toEqual([]) + expect(applyOperatorCassetteOps([op, op]).applied).toEqual([]) + expect(loadCassettes().find((c) => c.position === 1)!.count).toBe(80) + }) + + it('applies only the unseen ops from a window that mixes both', () => { + applyOperatorCassetteOps([ + { id: 'op-1', at: 1_700_000_000, type: 'refill', position: 1, bills: 10 }, + ]) + const result = applyOperatorCassetteOps([ + { id: 'op-1', at: 1_700_000_000, type: 'refill', position: 1, bills: 10 }, + { id: 'op-2', at: 1_700_000_001, type: 'refill', position: 1, bills: 5 }, + ]) + expect(result.applied).toEqual(['op-2']) + expect(loadCassettes().find((c) => c.position === 1)!.count).toBe(65) + }) + + it('applies a window oldest-first regardless of arrival order', () => { + // A recount then a refill is not the same as the reverse, so ordering is + // load-bearing and cannot be left to however the array arrived. + applyOperatorCassetteOps([ + { id: 'op-b', at: 1_700_000_002, type: 'refill', position: 1, bills: 7 }, + { id: 'op-a', at: 1_700_000_001, type: 'recount', position: 1, count: 3 }, + ]) + expect(loadCassettes().find((c) => c.position === 1)!.count).toBe(10) + }) + + it('empties a bay and sets a denomination', () => { + applyOperatorCassetteOps([ + { id: 'op-1', at: 1_700_000_000, type: 'empty', position: 2 }, + { id: 'op-2', at: 1_700_000_001, type: 'set_denomination', position: 2, denomination: 10 }, + ]) + const bay = loadCassettes().find((c) => c.position === 2)! + expect(bay.count).toBe(0) + expect(bay.denomination).toBe(10) + }) + + it('rejects a malformed op without applying or recording it', () => { + // Unrecorded on purpose: it stays pending on the operator's dashboard, + // which is the honest outcome. Recording it as applied would silence the + // noise by telling the operator their refill landed. + const result = applyOperatorCassetteOps([ + { id: 'op-1', at: 1_700_000_000, type: 'refill', position: 1, bills: -5 }, + ]) + expect(result.applied).toEqual([]) + expect(result.rejected[0]!.id).toBe('op-1') + expect(loadCassettes().find((c) => c.position === 1)!.count).toBe(50) + expect(getAppliedOpIds()).not.toContain('op-1') + }) + + it('echoes applied ids back, newest first', () => { + applyOperatorCassetteOps([ + { id: 'op-1', at: 1_700_000_000, type: 'refill', position: 1, bills: 1 }, + { id: 'op-2', at: 1_700_000_001, type: 'refill', position: 1, bills: 1 }, + ]) + expect(getAppliedOpIds()).toContain('op-1') + expect(getAppliedOpIds()).toContain('op-2') + }) + + it('advances the sequence on an applied op but not on a duplicate', () => { + const op = { + id: 'op-1', + at: 1_700_000_000, + type: 'refill' as const, + position: 1, + bills: 1, + } + const before = getCassetteStateSeq() + applyOperatorCassetteOps([op]) + const after = getCassetteStateSeq() + expect(after).toBeGreaterThan(before) + applyOperatorCassetteOps([op]) + expect(getCassetteStateSeq()).toBe(after) + }) + + it('advances the sequence on a dispense', () => { + const before = getCassetteStateSeq() + recordTransaction({ + ...TX_BASE, + txid: 'tx-seq', + type: 'cash_out', + status: 'complete', + bills: [{ denomination: 20, count: 1 }], + cassettes: [ + { + name: 'cassette1', + position: 1, + denomination: 20, + provisioned: 1, + dispensed: 1, + rejected: 0, + }, + ], + }) + expect(getCassetteStateSeq()).toBeGreaterThan(before) + }) +}) diff --git a/apps/machine/electron/main.ts b/apps/machine/electron/main.ts index 7383552..b5801ad 100644 --- a/apps/machine/electron/main.ts +++ b/apps/machine/electron/main.ts @@ -30,14 +30,17 @@ import { markStatePublished, resetStatePublishWatermark, resetForRepair, - applyOperatorCassettesConfig, + applyOperatorCassetteOps, + getAppliedOpIds, + getCassetteStateSeq, getFeeConfig, getLastKnownFeeConfigCreatedAt, applyFeeConfig, getBunkerBinding, saveBunkerBinding, clearBunkerBinding, - type OperatorCassettesPayload, + type CassetteOp, + type ApplyOpsResult, type FeeConfigPayload, type FeeConfigRow, type ApplyResult, @@ -566,10 +569,13 @@ ipcMain.handle('state:mark-state-published', (_event, unixTimestamp: number): vo markStatePublished(unixTimestamp) }) ipcMain.handle( - 'state:apply-operator-cassettes-config', - (_event, payload: OperatorCassettesPayload, eventCreatedAt: number): ApplyResult => - applyOperatorCassettesConfig(payload, eventCreatedAt) + 'state:apply-operator-cassette-ops', + (_event, ops: CassetteOp[]): ApplyOpsResult => applyOperatorCassetteOps(ops) ) +ipcMain.handle('state:get-applied-op-ids', (_event, limit?: number): string[] => + getAppliedOpIds(limit) +) +ipcMain.handle('state:get-cassette-state-seq', (): number => getCassetteStateSeq()) // Operator-fees consumer (aiolabs/lamassu-next#57) — persisted singleton // fee config + per-d-tag replay watermark + atomic apply for kind-30078 diff --git a/apps/machine/electron/preload.ts b/apps/machine/electron/preload.ts index c376170..8bebc13 100644 --- a/apps/machine/electron/preload.ts +++ b/apps/machine/electron/preload.ts @@ -163,13 +163,23 @@ contextBridge.exposeInMainWorld('electronAPI', { }): Promise<{ ok: boolean; bolt11?: string; reason?: string }> => ipcRenderer.invoke('lnurl:pay-session', args), - applyOperatorCassettesConfig: ( - payload: { - positions: Record - }, - eventCreatedAt: number - ): Promise<{ applied: true } | { applied: false; reason: string }> => - ipcRenderer.invoke('state:apply-operator-cassettes-config', payload, eventCreatedAt), + applyOperatorCassetteOps: ( + ops: { + id: string + at: number + type: 'refill' | 'empty' | 'recount' | 'set_denomination' + position: number + bills?: number + count?: number + denomination?: number + }[] + ): Promise<{ + applied: string[] + rejected: { id: string; reason: string }[] + }> => ipcRenderer.invoke('state:apply-operator-cassette-ops', ops), + getAppliedOpIds: (limit?: number): Promise => + ipcRenderer.invoke('state:get-applied-op-ids', limit), + getCassetteStateSeq: (): Promise => ipcRenderer.invoke('state:get-cassette-state-seq'), // Operator-fees consumer (aiolabs/lamassu-next#57) getFeeConfig: (): Promise<{ @@ -314,10 +324,22 @@ declare global { lnurlw: string amountMsat: number }) => Promise<{ ok: boolean; bolt11?: string; reason?: string }> - applyOperatorCassettesConfig: ( - payload: { positions: Record }, - eventCreatedAt: number - ) => Promise<{ applied: true } | { applied: false; reason: string }> + applyOperatorCassetteOps: ( + ops: { + id: string + at: number + type: 'refill' | 'empty' | 'recount' | 'set_denomination' + position: number + bills?: number + count?: number + denomination?: number + }[] + ) => Promise<{ + applied: string[] + rejected: { id: string; reason: string }[] + }> + getAppliedOpIds: (limit?: number) => Promise + getCassetteStateSeq: () => Promise getFeeConfig: () => Promise<{ cashInFeeFraction: number cashOutFeeFraction: number diff --git a/apps/machine/electron/state-store.ts b/apps/machine/electron/state-store.ts index ab3e304..fa6f17f 100644 --- a/apps/machine/electron/state-store.ts +++ b/apps/machine/electron/state-store.ts @@ -15,7 +15,7 @@ import fs from 'node:fs' let db: Database.Database | null = null -const SCHEMA_VERSION = '12' +const SCHEMA_VERSION = '13' function getDbPath(): string { const prodDir = '/var/lib/bitspire' @@ -57,6 +57,17 @@ export function initDatabase(dbPath?: string): void { count INTEGER NOT NULL DEFAULT 0 ); + CREATE TABLE IF NOT EXISTS cassette_ops ( + id TEXT PRIMARY KEY, + position INTEGER NOT NULL, + op_type TEXT NOT NULL, + bills INTEGER, + count INTEGER, + denomination INTEGER, + op_at INTEGER NOT NULL, + applied_at INTEGER NOT NULL + ); + CREATE TABLE IF NOT EXISTS cashbox ( id INTEGER PRIMARY KEY CHECK (id = 1), total_bills INTEGER NOT NULL DEFAULT 0, @@ -373,12 +384,47 @@ export function initDatabase(dbPath?: string): void { console.log('[StateStore] Migrated schema v11 → v12 (bunker_binding transport config)') } + if (existing && existing.value === '12') { + // Migration v12 → v13: operator OPERATIONS replace operator counts + // (aiolabs/bitspire ADR-004). + // + // The operator used to publish absolute counts and this machine applied + // them outright. Both sides wrote the same value over a transport that + // never tells a writer it lost, so a dashboard form loaded before a + // dispense silently discarded that dispense — and nothing on either side + // could detect it afterwards. The operator now publishes what it DID and + // this machine, which holds the notes, owns the running total. + // + // `cassette_ops` is the dedup ledger. A delta applied twice is wrong, and + // addressable events are re-delivered on every reconnect, so the operator + // mints an id per operation and we record the ones we have applied. The + // operator's window is a slice of recent operations rather than just the + // newest, so one we missed arrives with the next publish; dedup is what + // makes re-delivery free instead of dangerous. + db.exec(` + CREATE TABLE IF NOT EXISTS cassette_ops ( + id TEXT PRIMARY KEY, + position INTEGER NOT NULL, + op_type TEXT NOT NULL, + bills INTEGER, + count INTEGER, + denomination INTEGER, + op_at INTEGER NOT NULL, + applied_at INTEGER NOT NULL + ); + `) + db.prepare('UPDATE meta SET value = ? WHERE key = ?').run('13', 'schema_version') + console.log('[StateStore] Migrated schema v12 → v13 (added cassette_ops)') + existing.value = '13' + } + // 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 (?, ?)') seedMeta.run('lastKnownConfigCreatedAt', '0') seedMeta.run('bootstrapPublishedAt', '') seedMeta.run('lastKnownFeeConfigCreatedAt', '0') + seedMeta.run('cassetteStateSeq', '0') const cashboxRow = db.prepare('SELECT id FROM cashbox WHERE id = 1').get() if (!cashboxRow) { @@ -471,6 +517,36 @@ export function clearCountsUncertain(): void { ).run('countsUncertainSince', '') } +/** + * A counter bumped on every local change to a bay count, from any cause. + * + * It rides along in the state document so a reader can reject a regression + * without trusting a clock. `created_at` cannot carry that: it has + * second granularity, so two publishes in the same second are ordered by + * whichever event id hashes lower — and a machine whose clock stepped + * backwards would otherwise have every later report look older than the one + * already on the relay. + */ +export function getCassetteStateSeq(): number { + if (!db) throw new Error('Database not initialized') + const row = db.prepare('SELECT value FROM meta WHERE key = ?').get('cassetteStateSeq') as + | { value: string } + | undefined + return row ? Number(row.value) || 0 : 0 +} + +/** + * Bump the counter. Safe to call inside an open transaction — every caller + * that mutates a count does, so the bump commits or rolls back with it. + */ +export function bumpCassetteStateSeq(): void { + if (!db) throw new Error('Database not initialized') + db.prepare( + 'INSERT INTO meta (key, value) VALUES (?, ?) ' + + 'ON CONFLICT(key) DO UPDATE SET value = CAST(CAST(meta.value AS INTEGER) + 1 AS TEXT)' + ).run('cassetteStateSeq', '1') +} + /** Record the `created_at` just published, as the next publish's floor. */ export function markStatePublished(unixTimestamp: number): void { if (!db) throw new Error('Database not initialized') @@ -617,111 +693,190 @@ export function resetForRepair(): void { })() } -export type OperatorCassettesPayload = { - positions: Record -} - +/** + * Outcome of applying an operator-authored absolute config. Still the right + * shape for fee config, where the operator is the only writer of the value + * and a later event simply supersedes an earlier one. Cassette counts left + * this model in ADR-004 precisely because they had two writers. + */ export type ApplyResult = { applied: true } | { applied: false; reason: string } +/** One operator-authored operation, as it arrives on the wire. */ +export type CassetteOp = { + id: string + at: number + type: 'refill' | 'empty' | 'recount' | 'set_denomination' + position: number + bills?: number + count?: number + denomination?: number +} + +export type ApplyOpsResult = { + /** Ids applied by this call. Empty when every op was already on file. */ + applied: string[] + /** Ids rejected, with why. These stay unapplied and unrecorded. */ + rejected: { id: string; reason: string }[] +} + +const CASSETTE_OP_TYPES = new Set(['refill', 'empty', 'recount', 'set_denomination']) + /** - * Atomic apply of an operator-published cassette config (aiolabs/lamassu-next#56). + * Validate one operation in isolation. Returns null when it is well-formed. * - * Caller has already verified the event signature and decrypted the - * content. This function: - * - * 1. Rechecks replay-protection against `meta.lastKnownConfigCreatedAt` - * (defense-in-depth — caller should have done this too). - * 2. Validates the payload's `positions` key set is *exactly* the set of - * positions currently in the `cassettes` table. The bay count is - * hardware-determined and can't be added to or removed from via this - * path; only the per-bay denomination and count are operator-mutable. - * 3. Validates per-entry `denomination` is a positive int, `count` is a - * non-negative int. **Duplicate denominations across positions are - * intentionally permitted** — real machines load multiple cassettes - * with the same denomination for cash-out throughput. - * 4. In a single SQLite transaction: updates `cassettes` rows by position - * (denomination + count both mutable per row) AND advances - * `meta.lastKnownConfigCreatedAt` to `eventCreatedAt`. - * - * Mid-write crashes roll back cleanly; on restart the same event is - * re-delivered by the relay and the watermark check drops it as already - * consumed (or the watermark is pre-event because the tx rolled back, - * and the apply runs again from scratch). + * Shape errors and unknown positions are treated the same way by the caller: + * the op is neither applied nor recorded, so it stays pending on the + * operator's dashboard. That is the honest outcome — it did not happen — and + * it beats recording it as applied to stop the noise, which would tell the + * operator their refill landed when the notes are unaccounted for. */ -export function applyOperatorCassettesConfig( - payload: OperatorCassettesPayload, - eventCreatedAt: number -): ApplyResult { +function validateCassetteOp(op: CassetteOp, knownPositions: Set): string | null { + if (typeof op.id !== 'string' || op.id.length === 0) return 'missing id' + if (!CASSETTE_OP_TYPES.has(op.type)) return `unknown type ${String(op.type)}` + if (!Number.isInteger(op.position)) return `position must be an integer (got ${op.position})` + if (!knownPositions.has(op.position)) return `unknown position ${op.position}` + if (!Number.isFinite(op.at)) return 'missing at' + + if (op.type === 'refill') { + if (!Number.isInteger(op.bills) || (op.bills as number) <= 0) { + return `refill needs a positive integer bills (got ${op.bills})` + } + } + if (op.type === 'recount') { + if (!Number.isInteger(op.count) || (op.count as number) < 0) { + return `recount needs a non-negative integer count (got ${op.count})` + } + } + if (op.type === 'set_denomination') { + if (!Number.isInteger(op.denomination) || (op.denomination as number) <= 0) { + return `set_denomination needs a positive integer denomination (got ${op.denomination})` + } + } + return null +} + +/** + * Apply an operator's cassette operations, skipping any already on file. + * + * This replaces applying absolute counts. The operator authors what it DID — + * a refill in notes added, an empty, a recount, a denomination change — and + * this machine, which holds the physical notes, keeps the running total. + * Nobody but this process writes a count any more, so there is no second + * writer to lose a race to. + * + * Deltas are not idempotent and addressable events ARE re-delivered on every + * relay reconnect, so idempotency is carried explicitly: the operator mints an + * id per operation, `cassette_ops` records the ones applied, and a repeat is a + * no-op. That is also why there is no `created_at` watermark here any more. + * Under absolute counts the watermark was the only replay defence; with + * per-op ids it is strictly weaker than the dedup and would do active harm, + * because an event that arrives out of order may still carry an operation this + * machine has never seen. + * + * Applied oldest-first by `at`, ties broken by id so two operations stamped in + * the same second still order the same way on every machine. Ordering matters + * because a recount followed by a refill is not the same as the reverse. + * + * The whole batch runs in one SQLite transaction with the sequence bump, so a + * crash mid-apply rolls back to a coherent count and the next publish re-offers + * every op in the window. + */ +export function applyOperatorCassetteOps(ops: CassetteOp[]): ApplyOpsResult { if (!db) throw new Error('Database not initialized') + const database = db + const result: ApplyOpsResult = { applied: [], rejected: [] } + if (ops.length === 0) return result - const watermark = getLastKnownConfigCreatedAt() - if (eventCreatedAt <= watermark) { - return { - applied: false, - reason: `event.created_at (${eventCreatedAt}) <= lastKnownConfigCreatedAt (${watermark})`, - } - } - - const currentRows = db.prepare('SELECT position FROM cassettes').all() as { position: number }[] - const currentPositions = new Set(currentRows.map((r) => r.position)) - const payloadPositions = new Set(Object.keys(payload.positions).map((k) => Number(k))) - - if (currentPositions.size !== payloadPositions.size) { - return { - applied: false, - reason: `position count mismatch: state.db has ${currentPositions.size}, payload has ${payloadPositions.size}`, - } - } - for (const p of currentPositions) { - if (!payloadPositions.has(p)) { - return { applied: false, reason: `payload missing position ${p}` } - } - } - for (const p of payloadPositions) { - if (!currentPositions.has(p)) { - return { applied: false, reason: `payload includes unknown position ${p}` } - } - } - - for (const [posKey, entry] of Object.entries(payload.positions)) { - if (!Number.isInteger(entry.denomination) || entry.denomination <= 0) { - return { - applied: false, - reason: `denomination must be positive int (position ${posKey}, got ${entry.denomination})`, - } - } - if (!Number.isInteger(entry.count) || entry.count < 0) { - return { - applied: false, - reason: `count must be non-negative int (position ${posKey}, got ${entry.count})`, - } - } - } - - const updateCassette = db.prepare( - 'UPDATE cassettes SET denomination = ?, count = ? WHERE position = ?' + const knownPositions = new Set( + (database.prepare('SELECT position FROM cassettes').all() as { position: number }[]).map( + (r) => r.position + ) ) - const setWatermark = db.prepare('UPDATE meta SET value = ? WHERE key = ?') + const seen = database.prepare('SELECT 1 FROM cassette_ops WHERE id = ?') - const clearUncertain = db.prepare( + const pending: CassetteOp[] = [] + for (const op of ops) { + if (op && typeof op.id === 'string' && seen.get(op.id)) continue + const reason = validateCassetteOp(op, knownPositions) + if (reason) { + result.rejected.push({ id: op?.id ?? '', reason }) + continue + } + pending.push(op) + } + if (pending.length === 0) return result + + pending.sort((a, b) => a.at - b.at || (a.id < b.id ? -1 : a.id > b.id ? 1 : 0)) + + const addBills = database.prepare( + 'UPDATE cassettes SET count = MAX(0, count + ?) WHERE position = ?' + ) + const setCount = database.prepare('UPDATE cassettes SET count = ? WHERE position = ?') + const setDenomination = database.prepare( + 'UPDATE cassettes SET denomination = ? WHERE position = ?' + ) + const recordOp = database.prepare( + 'INSERT INTO cassette_ops (id, position, op_type, bills, count, denomination, op_at, applied_at) ' + + 'VALUES (?, ?, ?, ?, ?, ?, ?, ?)' + ) + const upsertMeta = database.prepare( 'INSERT INTO meta (key, value) VALUES (?, ?) ON CONFLICT(key) DO UPDATE SET value = excluded.value' ) - const run = db.transaction(() => { - for (const [posKey, entry] of Object.entries(payload.positions)) { - updateCassette.run(entry.denomination, entry.count, Number(posKey)) + const appliedAt = Math.floor(Date.now() / 1000) + let sawRecount = false + + database.transaction(() => { + for (const op of pending) { + if (op.type === 'refill') addBills.run(op.bills, op.position) + else if (op.type === 'empty') setCount.run(0, op.position) + else if (op.type === 'recount') { + setCount.run(op.count, op.position) + sawRecount = true + } else setDenomination.run(op.denomination, op.position) + + recordOp.run( + op.id, + op.position, + op.type, + op.bills ?? null, + op.count ?? null, + op.denomination ?? null, + Math.floor(op.at), + appliedAt + ) + result.applied.push(op.id) } - setWatermark.run(String(eventCreatedAt), 'lastKnownConfigCreatedAt') - // The operator just asserted real counts, which is what a recount is. - // Whatever made the old numbers untrustworthy no longer applies. - clearUncertain.run('countsUncertainSince', '') - }) + bumpCassetteStateSeq() + // A recount is an operator opening the bay and counting it, which is + // exactly what resolves an unverified count. Nothing else does: a refill + // adds to a number still known to be wrong. + if (sawRecount) upsertMeta.run('countsUncertainSince', '') + })() - run() console.log( - `[StateStore] Applied operator cassettes config @ created_at=${eventCreatedAt} (${Object.keys(payload.positions).length} positions)` + `[StateStore] Applied ${result.applied.length} cassette op(s)` + + (result.rejected.length ? `, rejected ${result.rejected.length}` : '') ) - return { applied: true } + return result +} + +/** + * The ids most recently applied, newest first — the acknowledgement leg of + * the protocol. + * + * An addressable event gives its publisher no failure signal at all: the relay + * returns OK for an event it then discards, and a losing writer is never told. + * Echoing the ids back in this machine's own state document is the only way + * the operator can distinguish an operation that landed from one that was + * merely sent. + */ +export function getAppliedOpIds(limit = 50): string[] { + if (!db) throw new Error('Database not initialized') + const rows = db + .prepare('SELECT id FROM cassette_ops ORDER BY applied_at DESC, rowid DESC LIMIT ?') + .all(limit) as { id: string }[] + return rows.map((r) => r.id) } // --------------------------------------------------------------------------- @@ -899,6 +1054,7 @@ export function setCassettes( const row = rows[i]! upsert.run(row.position ?? i + 1, row.denomination, row.count) } + bumpCassetteStateSeq() } ) @@ -913,11 +1069,14 @@ export function setCassettes( */ export function updateCassetteCountByPosition(position: number, delta: number): void { if (!db) throw new Error('Database not initialized') + const database = db - db.prepare('UPDATE cassettes SET count = MAX(0, count + ?) WHERE position = ?').run( - delta, - position - ) + database.transaction(() => { + database + .prepare('UPDATE cassettes SET count = MAX(0, count + ?) WHERE position = ?') + .run(delta, position) + bumpCassetteStateSeq() + })() } /** @@ -1120,6 +1279,10 @@ export function recordTransaction(tx: TransactionInput): void { } } } + // The counts moved, so the sequence must move with them, inside this + // same transaction. It rides in the state document as the operator's + // way to reject a regression without trusting either clock. + bumpCassetteStateSeq() } if (t.type === 'cash_in') { diff --git a/apps/machine/src/services/operator-config.ts b/apps/machine/src/services/operator-config.ts index 1627f5a..2f33f73 100644 --- a/apps/machine/src/services/operator-config.ts +++ b/apps/machine/src/services/operator-config.ts @@ -1,11 +1,17 @@ /** - * Operator-config consumer (aiolabs/lamassu-next#56). + * Operator-config consumer (aiolabs/lamassu-next#56, v2 per bitspire ADR-004). * * Subscribes to operator-published kind-30078 events carrying cassette - * config updates, validates + applies them to state.db, and hot-reloads - * the HAL dispenser. Also publishes a one-shot ATM-state hello-event on - * first boot so the operator dashboard (satmachineadmin) can auto-populate - * `cassette_configs` rows for this machine. + * OPERATIONS — a refill, an empty, a recount, a denomination change — + * applies the ones it has not already seen, and hot-reloads the HAL + * dispenser. It also publishes this machine's cassette state, which is + * what populates the operator dashboard's bay rows. + * + * The operator used to publish absolute counts and this machine applied them + * outright. Both sides wrote the same value over a transport that never tells + * a writer it lost: a dashboard form loaded before a dispense silently + * discarded that dispense, and neither side could detect it. This machine now + * owns the count — it holds the notes — and the operator says what it did. * * Architecture (see ~/dev/coordination/log.md entries on 2026-05-30): * @@ -35,6 +41,20 @@ import type {} from '@/types/electron' const KIND_NIP78 = 30078 +/** The wire schema this machine speaks. Operations, not counts (ADR-004). */ +const CASSETTE_SCHEMA_VERSION = 2 + +/** One operator-authored operation, as it arrives on the wire. */ +type CassetteOp = { + id: string + at: number + type: 'refill' | 'empty' | 'recount' | 'set_denomination' + position: number + bills?: number + count?: number + denomination?: number +} + /** Accept operator events stamped up to this many seconds in the future. */ const MAX_FUTURE_SKEW_S = 60 @@ -165,17 +185,13 @@ async function handleOperatorConfigEvent( return } - // 2. Replay protection — drop stale events. NIP-78 replaceable events - // DO get re-delivered on reconnect/restart; without this check, the - // ATM would re-apply the same payload on every boot and clobber any - // cash-out decrements that landed between operator publishes. - const watermark = await api.getLastKnownConfigCreatedAt() - if (event.created_at <= watermark) { - console.log( - `[OperatorConfig] Stale event dropped (created_at=${event.created_at} <= watermark=${watermark})` - ) - return - } + // 2. There is deliberately no `created_at` watermark here any more. + // Under absolute counts it was the only replay defence, and it cost us: + // an event re-delivered out of order was dropped whole, operations + // included. Idempotency now rides on the operations themselves — the + // operator mints an id per op and this machine records the ones it + // applied — which is strictly stronger, because it survives an event + // that mixes operations we have seen with ones we have not. // 3. Clock-skew defense — reject events stamped too far in the future. // Limits damage from a leaked operator nsec future-stamping a fake @@ -189,7 +205,7 @@ async function handleOperatorConfigEvent( } // 4. Decrypt content (NIP-44 v2). - let parsed: { positions: Record } + let parsed: { schema_version?: number; ops?: unknown } try { const plaintext = await cfg.signer.nip44Decrypt(event.pubkey, event.content) parsed = JSON.parse(plaintext) as typeof parsed @@ -197,22 +213,36 @@ async function handleOperatorConfigEvent( console.error('[OperatorConfig] Decrypt/parse failed:', err) return } - if (!parsed || typeof parsed !== 'object' || !parsed.positions) { - console.error('[OperatorConfig] Payload missing `positions` field') + if (!parsed || typeof parsed !== 'object' || !Array.isArray(parsed.ops)) { + // A v1 operator publishing absolute counts lands here and is ignored. + // That direction fails safe: the machine keeps its own counts, which it + // is now the only writer of, and simply will not dispense notes it + // believes it lacks. The opposite — applying a count from a form loaded + // before a dispense — is what ADR-004 exists to stop. + console.error('[OperatorConfig] Payload missing `ops` array — dropped') return } + const ops = parsed.ops as CassetteOp[] - // 5. Atomic apply (cassettes + meta watermark) via IPC. The state-store - // function re-validates watermark + position key-set equality + - // per-entry types inside the SQLite transaction. Duplicate - // denominations across positions are allowed — real machines load - // N cassettes of the same denomination for cash-out throughput. - const result = await api.applyOperatorCassettesConfig( - { positions: parsed.positions }, - event.created_at - ) - if (!result.applied) { - console.warn('[OperatorConfig] Apply rejected:', result.reason) + // 5. Apply the ones we have not seen, in one transaction with the sequence + // bump. No `created_at` watermark: each op carries an operator-minted id + // and the machine records what it applied, so a re-delivered event is a + // no-op on its own merits. The watermark would be strictly weaker and + // actively harmful — an event arriving out of order can still carry an + // operation this machine has never seen. + const result = await api.applyOperatorCassetteOps(ops) + for (const bad of result.rejected) { + console.warn(`[OperatorConfig] Op ${bad.id} rejected: ${bad.reason}`) + } + if (result.applied.length === 0) { + console.log(`[OperatorConfig] No new ops in event ${event.id.slice(0, 12)}…`) + // Still republish: the operator learns from our applied_ops echo that + // earlier operations landed, and an event carrying nothing new can be the + // first one we successfully answer after a relay outage. + const machineIdNoop = cfg.machineId ?? cfg.signer.pubkey + await publishCassettesState(cfg, api, machineIdNoop).catch((err) => + console.warn('[OperatorConfig] post-apply cassettes-state republish failed:', err) + ) return } @@ -220,8 +250,8 @@ async function handleOperatorConfigEvent( // picks up the new per-position mapping. state.db is already updated; // HAL re-init failure means the renderer's persistedInventory may be // ahead of the HAL until next service restart — log loudly but don't - // unwind the state.db apply (the operator wants their config landed; - // HAL can catch up). + // unwind the state.db apply (the operation happened physically; HAL + // can catch up). const cassettesAfter = await api.loadCassettes() const halResult = await api.halReloadCassettes( cassettesAfter.map((c) => ({ @@ -233,9 +263,7 @@ async function handleOperatorConfigEvent( if (!halResult.ok) { console.error('[OperatorConfig] HAL reload failed:', halResult.error) } - console.log( - `[OperatorConfig] Applied — created_at=${event.created_at}, positions=${Object.keys(parsed.positions).join(',')}` - ) + console.log(`[OperatorConfig] Applied ops: ${result.applied.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 @@ -274,9 +302,22 @@ async function publishCassettesState( // ignores this, so it needs no coordinated release. When set, the counts // above are the machine's best guess, not a measurement. const countsUncertainSince = await api.getCountsUncertainSince() - const payload = countsUncertainSince - ? { positions, counts_uncertain_since: countsUncertainSince } - : { positions } + + // `applied_ops` is the acknowledgement leg. An addressable event gives its + // publisher no failure signal — the relay returns OK for an event it then + // discards — so echoing the ids back is the only way the operator can tell + // an operation that landed from one that was merely sent. `seq` lets a + // reader reject a regression without trusting a clock: `created_at` has + // second granularity and ties break on event id, so it cannot order two + // reports from the same second. + const [appliedOps, seq] = await Promise.all([api.getAppliedOpIds(), api.getCassetteStateSeq()]) + const payload: Record = { + schema_version: CASSETTE_SCHEMA_VERSION, + positions, + seq, + applied_ops: appliedOps, + } + if (countsUncertainSince) payload.counts_uncertain_since = countsUncertainSince const ciphertext = await cfg.signer.nip44Encrypt(operatorPubkey, JSON.stringify(payload)) // Force the stamp strictly above our last one. Addressable events are ordered diff --git a/apps/machine/src/types/electron.d.ts b/apps/machine/src/types/electron.d.ts index e366f6d..8e4092e 100644 --- a/apps/machine/src/types/electron.d.ts +++ b/apps/machine/src/types/electron.d.ts @@ -183,10 +183,22 @@ declare global { pay: CardSession['pay'] amountMsat: number }) => Promise<{ ok: boolean; bolt11?: string; reason?: string }> - applyOperatorCassettesConfig: ( - payload: { positions: Record }, - eventCreatedAt: number - ) => Promise<{ applied: true } | { applied: false; reason: string }> + applyOperatorCassetteOps: ( + ops: { + id: string + at: number + type: 'refill' | 'empty' | 'recount' | 'set_denomination' + position: number + bills?: number + count?: number + denomination?: number + }[] + ) => Promise<{ + applied: string[] + rejected: { id: string; reason: string }[] + }> + getAppliedOpIds: (limit?: number) => Promise + getCassetteStateSeq: () => Promise getFeeConfig: () => Promise<{ cashInFeeFraction: number cashOutFeeFraction: number