Consume operator cassette operations instead of counts #106

Merged
padreug merged 2 commits from feat/cassette-ops-consumer into dev 2026-09-23 21:31:08 +00:00
7 changed files with 573 additions and 183 deletions

View file

@ -12,8 +12,10 @@
import { afterEach, beforeEach, describe, expect, it } from 'vitest' import { afterEach, beforeEach, describe, expect, it } from 'vitest'
import { import {
applyOperatorCassettesConfig, applyOperatorCassetteOps,
closeDatabase, closeDatabase,
getAppliedOpIds,
getCassetteStateSeq,
getCashbox, getCashbox,
getCountsUncertainSince, getCountsUncertainSince,
getInventory, getInventory,
@ -337,31 +339,157 @@ describe('state-store: unverified counts after a silent dispense', () => {
expect(getCountsUncertainSince()).toBe(1000) 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) markCountsUncertain(1000)
const applied = applyOperatorCassettesConfig( const result = applyOperatorCassetteOps([
{ { id: 'op-1', at: 1_700_000_000, type: 'recount', position: 1, count: 40 },
positions: { ])
'1': { denomination: 20, count: 40 }, expect(result.applied).toEqual(['op-1'])
'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.
expect(getCountsUncertainSince()).toBeNull() 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) markCountsUncertain(1000)
// Position key-set mismatch — the bay layout is hardware-determined. const result = applyOperatorCassetteOps([
const applied = applyOperatorCassettesConfig( { id: 'op-1', at: 1_700_000_000, type: 'refill', position: 1, bills: 10 },
{ positions: { '1': { denomination: 20, count: 40 } } }, ])
1_700_000_001 expect(result.applied).toEqual(['op-1'])
) expect(getCountsUncertainSince()).toBe(1000)
expect(applied.applied).toBe(false) })
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) 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)
})
})

View file

@ -30,14 +30,17 @@ import {
markStatePublished, markStatePublished,
resetStatePublishWatermark, resetStatePublishWatermark,
resetForRepair, resetForRepair,
applyOperatorCassettesConfig, applyOperatorCassetteOps,
getAppliedOpIds,
getCassetteStateSeq,
getFeeConfig, getFeeConfig,
getLastKnownFeeConfigCreatedAt, getLastKnownFeeConfigCreatedAt,
applyFeeConfig, applyFeeConfig,
getBunkerBinding, getBunkerBinding,
saveBunkerBinding, saveBunkerBinding,
clearBunkerBinding, clearBunkerBinding,
type OperatorCassettesPayload, type CassetteOp,
type ApplyOpsResult,
type FeeConfigPayload, type FeeConfigPayload,
type FeeConfigRow, type FeeConfigRow,
type ApplyResult, type ApplyResult,
@ -566,10 +569,13 @@ ipcMain.handle('state:mark-state-published', (_event, unixTimestamp: number): vo
markStatePublished(unixTimestamp) markStatePublished(unixTimestamp)
}) })
ipcMain.handle( ipcMain.handle(
'state:apply-operator-cassettes-config', 'state:apply-operator-cassette-ops',
(_event, payload: OperatorCassettesPayload, eventCreatedAt: number): ApplyResult => (_event, ops: CassetteOp[]): ApplyOpsResult => applyOperatorCassetteOps(ops)
applyOperatorCassettesConfig(payload, eventCreatedAt)
) )
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 // Operator-fees consumer (aiolabs/lamassu-next#57) — persisted singleton
// fee config + per-d-tag replay watermark + atomic apply for kind-30078 // fee config + per-d-tag replay watermark + atomic apply for kind-30078

View file

@ -163,13 +163,23 @@ contextBridge.exposeInMainWorld('electronAPI', {
}): Promise<{ ok: boolean; bolt11?: string; reason?: string }> => }): Promise<{ ok: boolean; bolt11?: string; reason?: string }> =>
ipcRenderer.invoke('lnurl:pay-session', args), ipcRenderer.invoke('lnurl:pay-session', args),
applyOperatorCassettesConfig: ( applyOperatorCassetteOps: (
payload: { ops: {
positions: Record<string, { denomination: number; count: number }> id: string
}, at: number
eventCreatedAt: number type: 'refill' | 'empty' | 'recount' | 'set_denomination'
): Promise<{ applied: true } | { applied: false; reason: string }> => position: number
ipcRenderer.invoke('state:apply-operator-cassettes-config', payload, eventCreatedAt), 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<string[]> =>
ipcRenderer.invoke('state:get-applied-op-ids', limit),
getCassetteStateSeq: (): Promise<number> => ipcRenderer.invoke('state:get-cassette-state-seq'),
// Operator-fees consumer (aiolabs/lamassu-next#57) // Operator-fees consumer (aiolabs/lamassu-next#57)
getFeeConfig: (): Promise<{ getFeeConfig: (): Promise<{
@ -314,10 +324,22 @@ declare global {
lnurlw: string lnurlw: string
amountMsat: number amountMsat: number
}) => Promise<{ ok: boolean; bolt11?: string; reason?: string }> }) => Promise<{ ok: boolean; bolt11?: string; reason?: string }>
applyOperatorCassettesConfig: ( applyOperatorCassetteOps: (
payload: { positions: Record<string, { denomination: number; count: number }> }, ops: {
eventCreatedAt: number id: string
) => Promise<{ applied: true } | { applied: false; reason: 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<string[]>
getCassetteStateSeq: () => Promise<number>
getFeeConfig: () => Promise<{ getFeeConfig: () => Promise<{
cashInFeeFraction: number cashInFeeFraction: number
cashOutFeeFraction: number cashOutFeeFraction: number

View file

@ -15,7 +15,7 @@ import fs from 'node:fs'
let db: Database.Database | null = null let db: Database.Database | null = null
const SCHEMA_VERSION = '12' const SCHEMA_VERSION = '13'
function getDbPath(): string { function getDbPath(): string {
const prodDir = '/var/lib/bitspire' const prodDir = '/var/lib/bitspire'
@ -57,6 +57,17 @@ export function initDatabase(dbPath?: string): void {
count INTEGER NOT NULL DEFAULT 0 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 ( CREATE TABLE IF NOT EXISTS cashbox (
id INTEGER PRIMARY KEY CHECK (id = 1), id INTEGER PRIMARY KEY CHECK (id = 1),
total_bills INTEGER NOT NULL DEFAULT 0, 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)') 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. // Defensive: a fresh install at SCHEMA_VERSION skips all migrations.
// Seed the operator-config meta rows if they're missing (idempotent). // Seed the operator-config meta rows if they're missing (idempotent).
const seedMeta = db.prepare('INSERT OR IGNORE INTO meta (key, value) VALUES (?, ?)') const seedMeta = db.prepare('INSERT OR IGNORE INTO meta (key, value) VALUES (?, ?)')
seedMeta.run('lastKnownConfigCreatedAt', '0') seedMeta.run('lastKnownConfigCreatedAt', '0')
seedMeta.run('bootstrapPublishedAt', '') seedMeta.run('bootstrapPublishedAt', '')
seedMeta.run('lastKnownFeeConfigCreatedAt', '0') seedMeta.run('lastKnownFeeConfigCreatedAt', '0')
seedMeta.run('cassetteStateSeq', '0')
const cashboxRow = db.prepare('SELECT id FROM cashbox WHERE id = 1').get() const cashboxRow = db.prepare('SELECT id FROM cashbox WHERE id = 1').get()
if (!cashboxRow) { if (!cashboxRow) {
@ -471,6 +517,36 @@ export function clearCountsUncertain(): void {
).run('countsUncertainSince', '') ).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. */ /** Record the `created_at` just published, as the next publish's floor. */
export function markStatePublished(unixTimestamp: number): void { export function markStatePublished(unixTimestamp: number): void {
if (!db) throw new Error('Database not initialized') if (!db) throw new Error('Database not initialized')
@ -617,111 +693,190 @@ export function resetForRepair(): void {
})() })()
} }
export type OperatorCassettesPayload = { /**
positions: Record<string, { denomination: number; count: number }> * 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 } 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 * Shape errors and unknown positions are treated the same way by the caller:
* content. This function: * 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
* 1. Rechecks replay-protection against `meta.lastKnownConfigCreatedAt` * it beats recording it as applied to stop the noise, which would tell the
* (defense-in-depth — caller should have done this too). * operator their refill landed when the notes are unaccounted for.
* 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).
*/ */
export function applyOperatorCassettesConfig( function validateCassetteOp(op: CassetteOp, knownPositions: Set<number>): string | null {
payload: OperatorCassettesPayload, if (typeof op.id !== 'string' || op.id.length === 0) return 'missing id'
eventCreatedAt: number if (!CASSETTE_OP_TYPES.has(op.type)) return `unknown type ${String(op.type)}`
): ApplyResult { 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') 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() const knownPositions = new Set(
if (eventCreatedAt <= watermark) { (database.prepare('SELECT position FROM cassettes').all() as { position: number }[]).map(
return { (r) => r.position
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 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 ?? '<no 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' 'INSERT INTO meta (key, value) VALUES (?, ?) ON CONFLICT(key) DO UPDATE SET value = excluded.value'
) )
const run = db.transaction(() => { const appliedAt = Math.floor(Date.now() / 1000)
for (const [posKey, entry] of Object.entries(payload.positions)) { let sawRecount = false
updateCassette.run(entry.denomination, entry.count, Number(posKey))
}
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', '')
})
run() database.transaction(() => {
console.log( for (const op of pending) {
`[StateStore] Applied operator cassettes config @ created_at=${eventCreatedAt} (${Object.keys(payload.positions).length} positions)` 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
) )
return { applied: true } result.applied.push(op.id)
}
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', '')
})()
console.log(
`[StateStore] Applied ${result.applied.length} cassette op(s)` +
(result.rejected.length ? `, rejected ${result.rejected.length}` : '')
)
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]! const row = rows[i]!
upsert.run(row.position ?? i + 1, row.denomination, row.count) 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 { export function updateCassetteCountByPosition(position: number, delta: number): void {
if (!db) throw new Error('Database not initialized') if (!db) throw new Error('Database not initialized')
const database = db
db.prepare('UPDATE cassettes SET count = MAX(0, count + ?) WHERE position = ?').run( database.transaction(() => {
delta, database
position .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') { if (t.type === 'cash_in') {

View file

@ -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 * Subscribes to operator-published kind-30078 events carrying cassette
* config updates, validates + applies them to state.db, and hot-reloads * OPERATIONS — a refill, an empty, a recount, a denomination change —
* the HAL dispenser. Also publishes a one-shot ATM-state hello-event on * applies the ones it has not already seen, and hot-reloads the HAL
* first boot so the operator dashboard (satmachineadmin) can auto-populate * dispenser. It also publishes this machine's cassette state, which is
* `cassette_configs` rows for this machine. * 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): * Architecture (see ~/dev/coordination/log.md entries on 2026-05-30):
* *
@ -35,6 +41,20 @@ import type {} from '@/types/electron'
const KIND_NIP78 = 30078 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. */ /** Accept operator events stamped up to this many seconds in the future. */
const MAX_FUTURE_SKEW_S = 60 const MAX_FUTURE_SKEW_S = 60
@ -165,17 +185,13 @@ async function handleOperatorConfigEvent(
return return
} }
// 2. Replay protection — drop stale events. NIP-78 replaceable events // 2. There is deliberately no `created_at` watermark here any more.
// DO get re-delivered on reconnect/restart; without this check, the // Under absolute counts it was the only replay defence, and it cost us:
// ATM would re-apply the same payload on every boot and clobber any // an event re-delivered out of order was dropped whole, operations
// cash-out decrements that landed between operator publishes. // included. Idempotency now rides on the operations themselves — the
const watermark = await api.getLastKnownConfigCreatedAt() // operator mints an id per op and this machine records the ones it
if (event.created_at <= watermark) { // applied — which is strictly stronger, because it survives an event
console.log( // that mixes operations we have seen with ones we have not.
`[OperatorConfig] Stale event dropped (created_at=${event.created_at} <= watermark=${watermark})`
)
return
}
// 3. Clock-skew defense — reject events stamped too far in the future. // 3. Clock-skew defense — reject events stamped too far in the future.
// Limits damage from a leaked operator nsec future-stamping a fake // Limits damage from a leaked operator nsec future-stamping a fake
@ -189,7 +205,7 @@ async function handleOperatorConfigEvent(
} }
// 4. Decrypt content (NIP-44 v2). // 4. Decrypt content (NIP-44 v2).
let parsed: { positions: Record<string, { denomination: number; count: number }> } let parsed: { schema_version?: number; ops?: unknown }
try { try {
const plaintext = await cfg.signer.nip44Decrypt(event.pubkey, event.content) const plaintext = await cfg.signer.nip44Decrypt(event.pubkey, event.content)
parsed = JSON.parse(plaintext) as typeof parsed parsed = JSON.parse(plaintext) as typeof parsed
@ -197,22 +213,36 @@ async function handleOperatorConfigEvent(
console.error('[OperatorConfig] Decrypt/parse failed:', err) console.error('[OperatorConfig] Decrypt/parse failed:', err)
return return
} }
if (!parsed || typeof parsed !== 'object' || !parsed.positions) { if (!parsed || typeof parsed !== 'object' || !Array.isArray(parsed.ops)) {
console.error('[OperatorConfig] Payload missing `positions` field') // 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 return
} }
const ops = parsed.ops as CassetteOp[]
// 5. Atomic apply (cassettes + meta watermark) via IPC. The state-store // 5. Apply the ones we have not seen, in one transaction with the sequence
// function re-validates watermark + position key-set equality + // bump. No `created_at` watermark: each op carries an operator-minted id
// per-entry types inside the SQLite transaction. Duplicate // and the machine records what it applied, so a re-delivered event is a
// denominations across positions are allowed — real machines load // no-op on its own merits. The watermark would be strictly weaker and
// N cassettes of the same denomination for cash-out throughput. // actively harmful — an event arriving out of order can still carry an
const result = await api.applyOperatorCassettesConfig( // operation this machine has never seen.
{ positions: parsed.positions }, const result = await api.applyOperatorCassetteOps(ops)
event.created_at 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)
) )
if (!result.applied) {
console.warn('[OperatorConfig] Apply rejected:', result.reason)
return return
} }
@ -220,8 +250,8 @@ async function handleOperatorConfigEvent(
// picks up the new per-position mapping. state.db is already updated; // picks up the new per-position mapping. state.db is already updated;
// HAL re-init failure means the renderer's persistedInventory may be // HAL re-init failure means the renderer's persistedInventory may be
// ahead of the HAL until next service restart — log loudly but don't // ahead of the HAL until next service restart — log loudly but don't
// unwind the state.db apply (the operator wants their config landed; // unwind the state.db apply (the operation happened physically; HAL
// HAL can catch up). // can catch up).
const cassettesAfter = await api.loadCassettes() const cassettesAfter = await api.loadCassettes()
const halResult = await api.halReloadCassettes( const halResult = await api.halReloadCassettes(
cassettesAfter.map((c) => ({ cassettesAfter.map((c) => ({
@ -233,9 +263,7 @@ async function handleOperatorConfigEvent(
if (!halResult.ok) { if (!halResult.ok) {
console.error('[OperatorConfig] HAL reload failed:', halResult.error) console.error('[OperatorConfig] HAL reload failed:', halResult.error)
} }
console.log( console.log(`[OperatorConfig] Applied ops: ${result.applied.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 // Republish our resulting cassette state so the operator's view reflects the
// applied config (the "on cassette reload" case). Different d-tag from 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 // ignores this, so it needs no coordinated release. When set, the counts
// above are the machine's best guess, not a measurement. // above are the machine's best guess, not a measurement.
const countsUncertainSince = await api.getCountsUncertainSince() const countsUncertainSince = await api.getCountsUncertainSince()
const payload = countsUncertainSince
? { positions, counts_uncertain_since: countsUncertainSince } // `applied_ops` is the acknowledgement leg. An addressable event gives its
: { positions } // 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<string, unknown> = {
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)) const ciphertext = await cfg.signer.nip44Encrypt(operatorPubkey, JSON.stringify(payload))
// Force the stamp strictly above our last one. Addressable events are ordered // Force the stamp strictly above our last one. Addressable events are ordered

View file

@ -183,10 +183,22 @@ declare global {
pay: CardSession['pay'] pay: CardSession['pay']
amountMsat: number amountMsat: number
}) => Promise<{ ok: boolean; bolt11?: string; reason?: string }> }) => Promise<{ ok: boolean; bolt11?: string; reason?: string }>
applyOperatorCassettesConfig: ( applyOperatorCassetteOps: (
payload: { positions: Record<string, { denomination: number; count: number }> }, ops: {
eventCreatedAt: number id: string
) => Promise<{ applied: true } | { applied: false; reason: 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<string[]>
getCassetteStateSeq: () => Promise<number>
getFeeConfig: () => Promise<{ getFeeConfig: () => Promise<{
cashInFeeFraction: number cashInFeeFraction: number
cashOutFeeFraction: number cashOutFeeFraction: number

View file

@ -65,21 +65,23 @@ The state document carries `applied_ops`, so the dashboard can render each publi
as applied or pending. This supplies the feedback leg a replaceable event cannot, without as applied or pending. This supplies the feedback leg a replaceable event cannot, without
needing the transport to report failures. needing the transport to report failures.
### 4a. Until decisions 1 to 4 ship, the overwrite is warned about, not prevented. ### 4a. Superseded. Before decisions 1 to 4 shipped, the overwrite was warned about.
The dashboard's publish dialog already states the failure plainly — that the publish will The dashboard's publish dialog stated the failure plainly — that the publish would overwrite
overwrite the ATM's tracked counts, that decrements since the last baseline will be lost, and the ATM's tracked counts, that decrements since the last baseline would be lost, and that it
that it should follow a physical refill rather than a mid-day tweak. It also says v2 should follow a physical refill rather than a mid-day tweak.
reconciliation will replace it.
Recording this because it changes how the gap should be read. It is a known, deliberately Kept here rather than deleted, because it is the calibration for how much a warning is worth.
accepted risk carrying a human-factors mitigation, not an oversight, and the product had It was a known, deliberately accepted risk carrying a human-factors mitigation, not an
already reached the same conclusion these decisions formalise. It is worth keeping in mind oversight, and the product had already reached the same conclusion these decisions formalise.
that a warning is the weakest control available: it depends on an operator reading a dialog It was also the weakest control available: it depended on an operator reading a dialog at the
at the end of a refill round, and it cannot help at all when the stale value is the one end of a refill round, and it could not help at all when the stale value was the one already
already in the form. Confirmed live on 2026-09-22 — a dispense moved a bay from 54 to 53 in the form. Confirmed live on 2026-09-22 — a dispense moved a bay from 54 to 53 while a form
while a form loaded at 54 stayed open, and nothing but that dialog stood between the operator loaded at 54 stayed open, and nothing but that dialog stood between the operator and
and discarding the decrement. discarding the decrement.
The dialog and the endpoint behind it are both gone. The operator dashboard no longer has a
field that accepts a count, which is a stronger guarantee than any wording could be.
### 5. Ordering is decided by `created_at`, never by arrival order, on both sides. ### 5. Ordering is decided by `created_at`, never by arrival order, on both sides.
@ -138,13 +140,29 @@ eight consecutive restarts, three heartbeat republishes carried strictly increas
back off the relay, and the machine, the relay and the operator dashboard agreed on the counts back off the relay, and the machine, the relay and the operator dashboard agreed on the counts
with timestamps correlated to the second. with timestamps correlated to the second.
Decisions 1 through 4 are the v2 operations wire and are not yet built. Decision 4a describes Decisions 1 through 4 are the v2 operations wire and shipped in spirekeeper#46 and
what stands in for them meanwhile. bitspire#106.
On the operator side there is no longer any endpoint that accepts a count: the absolute
publish, its CRUD write and its request model were removed rather than deprecated. The
dashboard records operations and renders each as applied or pending from the machine's
`applied_ops` echo. On the machine side, schema v13 adds a `cassette_ops` dedup ledger, the
`created_at` watermark on this path is retired in favour of per-op ids, and the state document
carries `schema_version`, `seq` and `applied_ops`.
The wire shapes are those given under decisions 1 and 4 above.
Cutover for v2 is strict, no compatibility code: spirekeeper deploys first, machines follow on Cutover for v2 is strict, no compatibility code: spirekeeper deploys first, machines follow on
their nightly pull. During that window a not-yet-updated ATM ignores an ops payload, so an their nightly pull. During that window a not-yet-updated ATM ignores an ops payload, so an
operator refill does not land until it updates — which fails safe, since the machine operator refill does not land until it updates — which fails safe, since the machine
under-counts and will not dispense bills it believes it lacks. under-counts and will not dispense bills it believes it lacks. In the other direction an
updated machine drops a v1 absolute-count payload on the missing `ops` array, which is the
same safe direction: the machine keeps the counts it is now the only writer of.
One gap stays open deliberately. An operation recorded while the relay is unreachable waits
for the operator's next action to be published, because only an operator action triggers a
publish. The window makes that self-healing once anything is published, but nothing on the
operator side republishes on its own. An operator-side heartbeat is the fix; it is not built.
## Alternatives considered ## Alternatives considered